Qdrant Edge를 서버와 동기화하기

Qdrant Edge를 서버와 동기화하기 (edge-edge-synchronization-guide)

Qdrant Edge는 외부 Qdrant 서버의 컬렉션과 동기화되어 다음과 같은 사용 사례를 지원할 수 있어요.

  • 인덱싱 오프로드(Offload indexing): 인덱싱은 계산 비용이 많이 드는 작업이에요. Edge 샤드를 서버 컬렉션과 동기화하면 인덱싱 과정을 더 강력한 서버 인스턴스로 오프로드할 수 있어요. 인덱싱된 데이터는 다시 Edge 샤드로 동기화될 수 있어요.
  • 백업 및 복원(Back up and Restore): Edge 샤드 데이터를 중앙 Qdrant 인스턴스에 정기적으로 백업해서 데이터 손실을 방지해요. 하드웨어 장애나 데이터 손상이 발생하면 중앙 인스턴스에서 데이터를 복원할 수 있어요.
  • 데이터 집계(Data Aggregation): 여러 위치에 배포된 여러 Edge 샤드에서 데이터를 수집하고, 포괄적인 분석과 리포팅을 위해 중앙 Qdrant 인스턴스로 집계해요.
  • 기기 간 동기화(Synchronization between devices): 여러 엣지 기기의 Edge 샤드를 중앙 Qdrant 인스턴스와 동기화해 데이터 일관성을 유지해요.

출처: Qdrant 공식 문서

Qdrant Edge를 서버와 동기화하기 (Synchronizing Qdrant Edge with a Server)

로컬 업데이트와 중앙 서버로부터의 업데이트를 모두 지원하려면, 두 개의 Edge 샤드로 구성된 셋업을 구현해요.

  • 로컬 데이터 업데이트를 처리하는 가변(mutable) Edge 샤드
  • 부분 스냅샷(partial snapshots)을 사용해 서버의 컬렉션에 있는 샤드를 미러링하는 불변(immutable) Edge 샤드

데이터를 쿼리할 때는 두 Edge 샤드의 결과를 병합해서 통합된 뷰를 제공해요. 이렇게 하면 로컬에서 새로 추가된 포인트가 서버에서 동기화된 데이터와 나란히 검색에 사용될 수 있어요.

Qdrant Edge 샤드는 중앙 서버와 동기화될 수 있어요

가변 Edge 샤드와 서버 컬렉션 둘 다에 데이터를 쓰는 이중 쓰기(dual-write) 메커니즘을 구현하면, 데이터가 서버에서 인덱싱되고 다시 불변 Edge 샤드로 동기화되어 검색 성능에 이점을 주게 돼요.

이 가이드에 설명된 패턴의 Python 예제 구현은 Qdrant Edge Demo GitHub 저장소에서 참고할 수 있어요.

1. 가변 Edge 샤드 초기화하기

가변 Edge 샤드는 로컬 데이터 업데이트를 관리해요. Qdrant Edge Quickstart 가이드에 설명된 대로 처음부터 초기화할 수 있어요.

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

MUTABLE_SHARD_DIR = "./qdrant-edge-directory/mutable"
Path(MUTABLE_SHARD_DIR).mkdir(parents=True, exist_ok=True)

VECTOR_NAME = "my-vector"
VECTOR_DIMENSION = 4

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

mutable_shard = EdgeShard.create(MUTABLE_SHARD_DIR, config)
const MUTABLE_SHARD_DIR: &str = "./qdrant-edge-directory/mutable";
const VECTOR_DIMENSION: usize = 4;
const VECTOR_NAME: &str = "my-vector";

fs_err::create_dir_all(MUTABLE_SHARD_DIR)?;

let config = EdgeConfig {
    on_disk_payload: true,
    vectors: HashMap::from([(
        VECTOR_NAME.to_string(),
        EdgeVectorParams {
            size: VECTOR_DIMENSION,
            distance: 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 mutable_shard = EdgeShard::new(Path::new(MUTABLE_SHARD_DIR), config)?;

2. 서버 스냅샷에서 불변 Edge 샤드 초기화하기

다음으로, 기존 Qdrant 컬렉션에서 Edge 샤드 초기화에서 설명한 대로 서버의 스냅샷으로 불변 Edge 샤드를 만듭니다.

import requests
import tempfile
import shutil

COLLECTION_NAME = "edge-collection"
snapshot_url = f"{QDRANT_URL}/collections/{COLLECTION_NAME}/shards/0/snapshot"
IMMUTABLE_SHARD_DIR = "./qdrant-edge-directory/immutable"

data_dir = Path(IMMUTABLE_SHARD_DIR)
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)

    immutable_shard = None
    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))

immutable_shard = EdgeShard.load(IMMUTABLE_SHARD_DIR)
const COLLECTION_NAME: &str = "edge-collection";
let snapshot_url = format!("{QDRANT_URL}/collections/{COLLECTION_NAME}/shards/0/snapshot");
const IMMUTABLE_SHARD_DIR: &str = "./qdrant-edge-directory/immutable";

let data_dir = Path::new(IMMUTABLE_SHARD_DIR);
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 immutable_shard = EdgeShard::load(data_dir, None)?;

3. 이중 쓰기 메커니즘 구현하기 (Implement a Dual-Write Mechanism)

두 Edge 샤드가 초기화됐다면, Edge 샤드에서 서버 컬렉션 업데이트에서 설명한 대로 애플리케이션에 이중 쓰기 메커니즘을 구현할 수 있어요.

먼저 서버에 써야 하는 대기 중인 업데이트를 담는 **큐(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
import time

SYNC_TIMESTAMP_KEY = "timestamp"

id = 2
vector = [0.4, 0.3, 0.2, 0.1]
payload = {"color": "green", SYNC_TIMESTAMP_KEY: time.time()}

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

rest_point = models.PointStruct(id=id, vector={VECTOR_NAME: vector}, payload=payload)
upload_queue.put(rest_point)
const SYNC_TIMESTAMP_KEY: &str = "timestamp";
let id = 2u64;
let vector = vec![0.4f32, 0.3, 0.2, 0.1];
let timestamp = SystemTime::now()
    .duration_since(UNIX_EPOCH)
    .unwrap()
    .as_secs_f64();
let payload = json!({"color": "green", SYNC_TIMESTAMP_KEY: timestamp});
let edge_points: Vec<PointStructPersisted> = vec![
    EdgePoint::new(
        PointId::NumId(id),
        Vectors::new_named([(VECTOR_NAME, vector.clone())]),
        payload.clone(),
    )
    .into(),
];
mutable_shard.update(UpdateOperation::PointOperation(
    PointOperations::UpsertPoints(PointInsertOperations::PointsList(edge_points)),
))?;

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

각 포인트의 payload에는 포인트가 업서트된 시점을 기록하는 타임스탬프 필드(이 예제에서는 SYNC_TIMESTAMP_KEY)가 포함돼야 해요. 이 타임스탬프는 불변 Edge 샤드가 서버와 동기화될 때 **데이터를 중복 제거(deduplicate)**하는 데 사용돼요.

백그라운드 워커가 업로드 큐를 처리하고 업데이트를 정기적으로 서버 컬렉션에 쓸 수 있어요. 이렇게 하면 서버 컬렉션이 가변 Edge 샤드의 로컬 변경 사항으로 최신 상태를 유지하는 걸 보장해요.

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?;
}

4. 불변 Edge 샤드 정기적으로 업데이트하기

Qdrant Edge를 서버측 변경으로 업데이트에서 설명한 대로, 부분 스냅샷을 사용해 불변 Edge 샤드를 서버의 변경 사항으로 주기적으로 업데이트할 수 있어요.

스냅샷을 복원하는 동안에는 가변 Edge 샤드의 진행 중인 데이터 업데이트를 일시 중지하고 버퍼링하고 싶을 수 있어요. 스냅샷을 만들기 전에 큐에 있는 모든 데이터가 서버에 기록됐는지 확인하세요. 복원이 완료된 뒤에는 정상 작업을 재개할 수 있어요. Python 예제 구현은 Qdrant Edge Demo GitHub 저장소를 참고하세요.

import time

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

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)

    immutable_shard.update_from_snapshot(str(partial_snapshot_path))
let sync_timestamp = SystemTime::now()
    .duration_since(UNIX_EPOCH)
    .unwrap()
    .as_secs_f64();
let current_manifest = immutable_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 immutable_shard = EdgeShard::recover_partial_snapshot(
    data_dir,
    &current_manifest,
    unpacked_dir.path(),
    &snapshot_manifest,
)?;

이 예제는 부분 스냅샷을 만들 때 sync_timestamp를 기록해요. 이 타임스탬프 이전에 가변 Edge 샤드에 추가된 모든 포인트가 이제 불변 Edge 샤드로 복원됐어요. 이 중복 포인트들은 이제 가변 Edge 샤드에서 삭제할 수 있어요.

from qdrant_edge import (
    Filter,
    FieldCondition,
    RangeFloat,
)

mutable_shard.update(
    UpdateOperation.delete_points_by_filter(
        Filter(
            must=[
                FieldCondition(
                    key=SYNC_TIMESTAMP_KEY,
                    range=RangeFloat(lte=sync_timestamp),
                )
            ]
        )
    )
)
let filter = Filter::new_must(Condition::Field(FieldCondition::new_range(
    SYNC_TIMESTAMP_KEY.parse::<JsonPath>().unwrap(),
    Range {
        lte: Some(OrderedFloat(sync_timestamp)),
        ..Default::default()
    },
)));
mutable_shard.update(UpdateOperation::PointOperation(
    PointOperations::DeletePointsByFilter(filter),
))?;

모든 데이터에 걸쳐 통합된 검색 경험을 제공하려면, 가변 Edge 샤드와 불변 Edge 샤드 둘 다를 쿼리하고 두 결과 집합을 병합해요. 포인트가 두 Edge 샤드에 모두 존재할 수 있으므로 포인트 ID 기준으로 결과를 중복 제거해요.

from qdrant_edge import Query, QueryRequest

query_request = QueryRequest(
    query=Query.Nearest([0.2, 0.1, 0.9, 0.7], using=VECTOR_NAME),
    limit=10,
    with_vector=False,
    with_payload=True,
)

mutable_results = mutable_shard.query(query_request)
immutable_results = immutable_shard.query(query_request)

all_results = list(mutable_results) + list(immutable_results)
all_results.sort(key=lambda x: x.score, reverse=True)

seen_ids = set()
unique_results = []
for result in all_results:
    if result.id not in seen_ids:
        seen_ids.add(result.id)
        unique_results.append(result)

results = [
    {"id": result.id, "score": result.score, "payload": result.payload}
    for result in unique_results[:10]
]
let query = QueryRequest {
    prefetches: vec![],
    query: Some(ScoringQuery::Vector(QueryEnum::Nearest(NamedQuery {
        query: vec![0.2f32, 0.1, 0.9, 0.7].into(),
        using: Some(VECTOR_NAME.to_string()),
    }))),
    filter: None,
    score_threshold: None,
    limit: 10,
    offset: 0,
    params: None,
    with_vector: WithVector::Bool(false),
    with_payload: WithPayloadInterface::Bool(true),
};

let mut all_results = mutable_shard.query(query.clone())?;
all_results.extend(immutable_shard.query(query)?);
all_results.sort_by(|a, b| {
    b.score.partial_cmp(&a.score).unwrap_or(std::cmp::Ordering::Equal)
});

let mut seen_ids = std::collections::HashSet::new();
let results: Vec<_> = all_results
    .into_iter()
    .filter(|p| seen_ids.insert(p.id.clone()))
    .take(10)
    .collect();

지원 (Support)

프로젝트에 Qdrant Edge를 구현하는 데 명시적 지원이 필요하면, Qdrant Sales에 문의하세요.

더 알아보기 (Learn more)