데이터 동기화 패턴

데이터 동기화 패턴 (edge-edge-data-synchronization-patterns)

이 페이지는 Qdrant Edge 샤드Qdrant 서버 컬렉션 사이에서 데이터를 동기화하는 패턴들을 설명해요. 이 패턴들을 실제로 구현하는 종단 간(end-to-end) 가이드는 Qdrant Edge 동기화 가이드를 참고하세요.

출처: Qdrant 공식 문서

기존 Qdrant 컬렉션에서 Edge 샤드 초기화하기 (Initialize Edge Shard from Existing Qdrant Collection)

비어 있는 Edge 샤드로 시작하는 대신, Qdrant 서버의 컬렉션에 있는 기존 데이터로 초기화하고 싶을 수 있어요. 서버측 컬렉션에 있는 샤드의 **스냅샷을 복원(restore)**해서 이 작업을 할 수 있어요.

Qdrant Edge 샤드는 서버측 샤드의 스냅샷에서 초기화될 수 있어요

동기화용 스냅샷을 만들 때 스냅샷 URL에 적용할 서버측 샤드 ID를 지정해요. 이렇게 하면 하나의 컬렉션이 여러 독립 사용자나 기기 — 각각 자신의 Edge 샤드를 가진 — 에게 서빙할 수 있어요. Qdrant의 샤딩 전략에 대해 더 읽으려면 티어형 멀티테넌시(Tiered Multitenancy) 문서를 참고하세요.

먼저 스냅샷 URL을 만들겠어요.

COLLECTION_NAME = "edge-collection"
snapshot_url = f"{QDRANT_URL}/collections/{COLLECTION_NAME}/shards/0/snapshot"
const COLLECTION_NAME: &str = "edge-collection";
let snapshot_url = format!("{QDRANT_URL}/collections/{COLLECTION_NAME}/shards/0/snapshot");

이 예제는 샤드 ID 0을 사용한다는 점을 참고하세요.

스냅샷 URL을 사용해서 스냅샷을 로컬 디스크에 다운로드하고, 그 데이터로 새 Edge 샤드를 초기화할 수 있어요.

from pathlib import Path
from qdrant_edge import EdgeShard
import requests
import shutil
import tempfile

SHARD_DIRECTORY = "./qdrant-edge-directory"
data_dir = Path(SHARD_DIRECTORY)

with tempfile.TemporaryDirectory(dir=data_dir.parent) as restore_dir:
    snapshot_path = Path(restore_dir) / "shard.snapshot"
    with requests.get(snapshot_url, headers={"api-key": QDRANT_API_KEY}, stream=True) as r:
        r.raise_for_status()
        with open(snapshot_path, "wb") as f:
            for chunk in r.iter_content(chunk_size=8192):
                f.write(chunk)

    if data_dir.exists():
        shutil.rmtree(data_dir)
    data_dir.mkdir(parents=True, exist_ok=True)

    EdgeShard.unpack_snapshot(str(snapshot_path), str(data_dir))

edge_shard = EdgeShard.load(SHARD_DIRECTORY)
const SHARD_DIRECTORY: &str = "./qdrant-edge-directory";
let data_dir = Path::new(SHARD_DIRECTORY);
let restore_dir = tempfile::Builder::new()
    .tempdir_in(data_dir.parent().unwrap_or(Path::new(".")))?;
let snapshot_path = restore_dir.path().join("shard.snapshot");

let mut bytes = Vec::new();
std::io::copy(
    &mut ureq::get(&snapshot_url)
        .header("api-key", QDRANT_API_KEY)
        .call()?
        .into_body()
        .into_reader(),
    &mut bytes,
)?;
fs_err::write(&snapshot_path, &bytes)?;

if data_dir.exists() {
    fs_err::remove_dir_all(data_dir)?;
}
fs_err::create_dir_all(data_dir)?;

EdgeShard::unpack_snapshot(&snapshot_path, data_dir)?;
let edge_shard = EdgeShard::load(data_dir, None)?;

이 코드는 먼저 스냅샷을 임시 디렉토리에 다운로드해요. 다음으로 EdgeShard.unpack_snapshot이 다운로드한 스냅샷을 데이터 디렉토리에 풀고, 언팩된 데이터와 구성으로 Edge 샤드를 초기화해요.

edge_shard는 스냅샷이 만들어진 소스 컬렉션과 같은 구성과 같은 파일 구조를 사용해요. 벡터 및 payload 인덱스도 포함해서요.

Qdrant Edge를 서버측 변경으로 업데이트하기 (Update Qdrant Edge with Server-Side Changes)

Edge 샤드를 서버 컬렉션의 새 데이터로 최신 상태로 유지하려면, 스냅샷을 주기적으로 다운로드해서 적용하면 돼요. 매번 전체 스냅샷을 복원하면 불필요한 오버헤드가 생겨요. 대신 **부분 스냅샷(partial snapshot)**을 사용해서 마지막 스냅샷 이후의 변경분만 복원할 수 있어요. 부분 스냅샷은 Edge 샤드의 모든 세그먼트와 메타데이터를 설명하는 매니페스트를 기준으로, 변경된 세그먼트만 담고 있어요.

manifest = edge_shard.snapshot_manifest()
url = f"{QDRANT_URL}/collections/{COLLECTION_NAME}/shards/0/snapshot/partial/create"

with tempfile.TemporaryDirectory(dir=data_dir) as temp_dir:
    partial_snapshot_path = Path(temp_dir) / "partial.snapshot"
    response = requests.post(url, headers={"api-key": QDRANT_API_KEY}, json=manifest, stream=True)
    response.raise_for_status()
    with open(partial_snapshot_path, "wb") as f:
        for chunk in response.iter_content(chunk_size=8192):
            f.write(chunk)

    edge_shard.update_from_snapshot(str(partial_snapshot_path))
let current_manifest = edge_shard.snapshot_manifest()?;
let update_url = format!("{QDRANT_URL}/collections/{COLLECTION_NAME}/shards/0/snapshot/partial/create");
let temp_dir = tempfile::tempdir_in(data_dir)?;
let partial_snapshot_path = temp_dir.path().join("partial.snapshot");

let mut bytes = Vec::new();
std::io::copy(
    &mut ureq::post(&update_url)
        .header("api-key", QDRANT_API_KEY)
        .send_json(&current_manifest)?
        .into_body()
        .into_reader(),
    &mut bytes,
)?;
fs_err::write(&partial_snapshot_path, &bytes)?;

let unpacked_dir = tempfile::tempdir_in(data_dir)?;
EdgeShard::unpack_snapshot(&partial_snapshot_path, unpacked_dir.path())?;
let snapshot_manifest = SnapshotManifest::load_from_snapshot(unpacked_dir.path(), None,)?;
let edge_shard = EdgeShard::recover_partial_snapshot(
    data_dir,
    &current_manifest,
    unpacked_dir.path(),
    &snapshot_manifest,
)?;

Edge 샤드에서 서버 컬렉션 업데이트하기 (Update a Server Collection from an Edge Shard)

Edge 샤드에서 서버 컬렉션으로 데이터를 동기화하려면, 애플리케이션에 이중 쓰기(dual-write) 메커니즘을 구현해야 해요. Edge 샤드에서 포인트를 추가하거나 업데이트할 때, 동시에 Qdrant 클라이언트로 서버 컬렉션에도 저장하면 돼요.

서버 컬렉션에 직접 쓰는 대신, 동기화를 비동기로 처리하는 백그라운드 작업이나 메시지 큐를 설정하고 싶을 수도 있어요. 애플리케이션이 항상 안정적인 인터넷 연결을 갖고 있지 않을 수 있으니, 업데이트를 큐에 넣어 두면 연결이 복구될 때 데이터가 결국 동기화되게 보장할 수 있어요.

먼저 다음을 초기화해요.

  • 처음부터 또는 서버측 스냅샷에서 만든 Edge 샤드
  • Qdrant 서버 연결

Edge 샤드 초기화:

from pathlib import Path
from qdrant_edge import (
    Distance,
    EdgeConfig,
    EdgeVectorParams,
)

SHARD_DIRECTORY = "./qdrant-edge-directory"
VECTOR_NAME = "my-vector"
VECTOR_DIMENSION = 4

Path(SHARD_DIRECTORY).mkdir(parents=True, exist_ok=True)

config = EdgeConfig(
    vectors={
        VECTOR_NAME: EdgeVectorParams(
            size=VECTOR_DIMENSION,
            distance=Distance.Cosine,
        )
    }
)

edge_shard = EdgeShard.create(SHARD_DIRECTORY, config)
const VECTOR_DIMENSION: usize = 4;
const VECTOR_NAME: &str = "my-vector";

fs_err::create_dir_all(SHARD_DIRECTORY)?;

let config = EdgeConfig {
    on_disk_payload: true,
    vectors: HashMap::from([(
        VECTOR_NAME.to_string(),
        EdgeVectorParams {
            size: VECTOR_DIMENSION,
            distance: qdrant_edge::Distance::Cosine,
            on_disk: Some(true),
            quantization_config: None,
            multivector_config: None,
            datatype: None,
            hnsw_config: None,
        },
    )]),
    sparse_vectors: HashMap::new(),
    hnsw_config: Default::default(),
    quantization_config: None,
    optimizers: Default::default(),
};

let edge_shard = EdgeShard::new(Path::new(SHARD_DIRECTORY), config)?;

서버에 대한 Qdrant 클라이언트 연결을 초기화하고, 대상 컬렉션이 없으면 만들 수 있어요.

from qdrant_client import QdrantClient, models

server_client = QdrantClient(url=QDRANT_URL, api_key=QDRANT_API_KEY)

COLLECTION_NAME = "edge-collection"
if not server_client.collection_exists(collection_name=COLLECTION_NAME):
    server_client.create_collection(
        collection_name=COLLECTION_NAME,
        vectors_config={
            VECTOR_NAME: models.VectorParams(size=VECTOR_DIMENSION, distance=models.Distance.COSINE)
        }
    )
let server_client = Qdrant::from_url(QDRANT_URL)
    .api_key(QDRANT_API_KEY)
    .build()?;
if !server_client.collection_exists(COLLECTION_NAME).await? {
    server_client
        .create_collection(
            CreateCollectionBuilder::new(COLLECTION_NAME)
                .vectors_config(VectorParamsBuilder::new(VECTOR_DIMENSION as u64, Distance::Cosine))
        )
        .await?;
}

다음으로, 서버와 동기화해야 하는 포인트들을 담는 **큐(queue)**를 만듭니다.

from queue import Empty, Queue
# This is in-memory queue
# For production use cases consider persisting changes
upload_queue: Queue[models.PointStruct] = Queue()
// This is an in-memory queue.
// For production use cases consider persisting changes.
let mut upload_queue: std::collections::VecDeque<PointStruct> = std::collections::VecDeque::new();

Edge 샤드에서 포인트를 추가하거나 업데이트할 때, 그 포인트를 서버와의 동기화용으로도 큐에 넣어요.

from qdrant_edge import (
    Point,
    UpdateOperation,
)
from qdrant_client import models

id = 1
vector = [0.1, 0.2, 0.3, 0.4]
payload = {"color": "red"}

point = Point(id=id, vector={VECTOR_NAME: vector}, payload=payload)
edge_shard.update(UpdateOperation.upsert_points([point]))

rest_point = models.PointStruct(id=id, vector={VECTOR_NAME: vector}, payload=payload)
upload_queue.put(rest_point)
let id = 1u64;
let vector = vec![0.1f32, 0.2, 0.3, 0.4];
let payload = json!({"color": "red"});
let edge_points: Vec<PointStructPersisted> = vec![
    EdgePoint::new(
        PointId::NumId(id),
        Vectors::new_named([(VECTOR_NAME, vector.clone())]),
        payload.clone(),
    )
    .into(),
];
edge_shard.update(UpdateOperation::PointOperation(
    PointOperations::UpsertPoints(PointInsertOperations::PointsList(edge_points)),
))?;

let server_point = PointStruct::new(
    id,
    HashMap::from([(VECTOR_NAME.to_string(), vector)]),
    payload.as_object().cloned().unwrap_or_default(),
);
upload_queue.push_back(server_point);

백그라운드 워커가 업로드 큐를 처리하고 포인트들을 서버 컬렉션과 동기화할 수 있어요. 이 예제는 한 번에 최대 10개 포인트씩 배치로 업로드해요.

BATCH_SIZE = 10

points_to_upload: list[models.PointStruct] = []
while len(points_to_upload) < BATCH_SIZE:
    try:
        points_to_upload.append(upload_queue.get_nowait())
    except Empty:
        break

if points_to_upload:
    server_client.upsert(
        collection_name=COLLECTION_NAME,
        points=points_to_upload,
    )
const BATCH_SIZE: usize = 10;
let points_to_upload: Vec<PointStruct> =
    upload_queue.drain(..BATCH_SIZE.min(upload_queue.len())).collect();
if !points_to_upload.is_empty() {
    server_client
        .upsert_points(UpsertPointsBuilder::new(COLLECTION_NAME, points_to_upload))
        .await?;
}

네트워크 문제나 서버 이용 불가 시 오류 처리와 재시도를 제대로 구현해야 한다는 점을 꼭 기억하세요.

더 알아보기 (Learn more)