CDC 복제에서의 대량 가져오기

이 가이드에서는 CDC 복제 토폴로지의 일부인 Milvus 클러스터에 대해 대량 가져오기를 실행하는 방법을 설명합니다. 복제 클러스터에서 대량 가져오기를 수행할 때는 2단계 커밋(2PC)을 사용해야 하며, 이를 통해 가져오기 작업이 주 클러스터와 대기 클러스터 전반에 걸쳐 단일하고 순서대로 정렬된 시점으로 커밋됩니다.

이 가이드에서 주 클러스터는 소스 Milvus 클러스터이며, 대기 클러스터는 대상 Milvus 클러스터입니다.

시작하기 전에 클러스터 간에 CDC 복제가 이미 구성되어 있는지 확인하십시오. 자세한 내용은 ‘CDC 복제 설정’을 참조하십시오.

2PC가 필요한 이유

일반적인 대량 가져오기는 가져오기 작업이 완료되면 자동으로 커밋되어, 가져온 데이터를 즉시 확인할 수 있습니다. 그러나 CDC 복제 토폴로지에서는 주 클러스터와 대기 클러스터가 가져온 데이터를 동일한 논리적 시점에서 노출해야 하므로 이러한 동작이 허용되지 않습니다.

대신, ` auto_commit=false`를 설정하여 2단계 커밋(two-phase commit) 모드로 가져오기를 실행하십시오:

  1. 가져오기 단계: Milvus는 프라이머리 클러스터에 데이터를 로드하고 가져오기 작업을 스탠바이 클러스터로 복제하지만, 가져온 데이터는 여전히 표시되지 않습니다. 가져오기 작업은 ‘ Uncommitted ’ 상태에서 중지되어 대기합니다.

  2. 커밋 단계: 프라이머리 클러스터에서 가져오기 작업을 명시적으로 커밋합니다. 커밋은 단일 순차적 펜스(fence)로 스탠바이 클러스터에 복제되므로, 두 클러스터 모두 동일한 논리적 시점에서 가져온 데이터를 표시하게 됩니다.

1단계: 복제 클러스터에서 가져오기 활성화

복제 클러스터에서의 가져오기 기능은 기본적으로 비활성화되어 있습니다. 주 클러스터와 대기 클러스터 모두에서 ` dataCoord.import.enableInReplicatingCluster `를 ` true `로 설정하여 이 기능을 활성화하십시오.

Milvus Operator를 사용하여 Milvus를 배포하는 경우, 각 Milvus 리소스의 spec.config 에 다음 설정을 추가하십시오:

spec:
  config:
    dataCoord:
      import:
        enableInReplicatingCluster: true

milvus.yaml 을 통해 Milvus를 직접 구성하는 경우 다음 설정을 추가하십시오:

dataCoord:
  import:
    enableInReplicatingCluster: true

이 설정은 새로 고침이 가능하므로, 클러스터를 완전히 재시작하지 않아도 적용될 수 있습니다.

이 설정이 활성화되면 복제 클러스터는 ` auto_commit=false`가 포함된 가져오기 요청만 수락합니다. 다음 표에는 일반적으로 거부되는 요청이 나열되어 있습니다:

상황오류 메시지
dataCoord.import.enableInReplicatingCluster 이 활성화되지 않은 경우import in replicating cluster is not supported yet
auto_commit=true 이 제출됨auto_commit=true import in replicating cluster is not supported

2단계: 2PC 가져오기 실행

모든 가져오기 호출을 주 클러스터에서 실행하십시오. 가져온 데이터와 커밋 결정은 스탠바이 클러스터로 자동으로 복제되므로, 스탠바이 클러스터에서 직접 가져오기를 제출하거나 커밋하지 마십시오.

각 클러스터는 자체 오브젝트 스토리지에서 가져오기 파일을 읽습니다. 가져올 파일이 주 클러스터와 대기 클러스터의 오브젝트 스토리지 모두에 존재하는지 확인하십시오. 파일을 두 클러스터 모두에 업로드하거나, 두 클러스터가 모두 읽을 수 있는 오브젝트 스토리지를 사용할 수 있습니다. 대기 클러스터에서 파일이 누락된 경우, 복제된 가져오기 작업은 ‘오브젝트 없음(object-not-found)’ 오류와 함께 실패합니다.

다음 예제에서는 pymilvus.bulk_writer 의 REST 기반 가져오기 헬퍼를 사용합니다. url 값은 다른 API 호출에 사용하는 것과 동일한 Milvus 주소입니다.

import time

from pymilvus.bulk_writer import bulk_import, commit_import, get_import_progress

primary_url = "http://127.0.0.1:19530"
standby_url = "http://127.0.0.1:19531"

collection_name = "demo_collection"

# Object-storage paths of the files to import. Prepare these files the same
# way as a normal bulk import, for example by using BulkWriter.
files = [
    ["import-data/part-1.parquet"],
]


def wait_for_state(url, job_id, target_state, timeout=600):
    deadline = time.time() + timeout
    while time.time() < deadline:
        resp = get_import_progress(url=url, job_id=job_id)
        data = resp.json().get("data", {})
        state = data.get("state")
        print(f"[{url}] job {job_id} state={state}, progress={data.get('progress')}")

        if state == target_state:
            return
        if state == "Failed":
            raise RuntimeError(
                f"import job {job_id} failed on {url}: {data.get('reason')}"
            )

        time.sleep(3)

    raise TimeoutError(f"job {job_id} did not reach {target_state} on {url}")


# Start a 2PC import on the primary cluster. In a replicating cluster,
# auto_commit=false is required, and the job stops at the Uncommitted state.
resp = bulk_import(
    url=primary_url,
    collection_name=collection_name,
    files=files,
    options={"auto_commit": "false"},
)
job_id = resp.json()["data"]["jobId"]
print(f"started 2PC import job: {job_id}")

# Wait until both clusters report Uncommitted. The same job ID is used on the
# primary and standby clusters because the import is replicated through CDC.
wait_for_state(primary_url, job_id, "Uncommitted")
wait_for_state(standby_url, job_id, "Uncommitted")

# Commit once on the primary cluster. Do not commit on the standby cluster.
commit_import(url=primary_url, job_id=job_id)
print(f"committed import job: {job_id}")

# Wait until the import is completed and visible on both clusters.
wait_for_state(primary_url, job_id, "Completed")
wait_for_state(standby_url, job_id, "Completed")
print("import committed and visible on both clusters")

두 클러스터 모두에서 Uncommitted 를 기다려야 하는 이유

대기 클러스터의 가져오기가 완료되기 전에 커밋을 수행해도 데이터가 손상되지는 않지만, 커밋이 적용될 때 대기 클러스터는 여전히 데이터를 따라잡는 중입니다. 주 클러스터와 대기 클러스터 모두 Uncommitted 를 반환할 때까지 기다리면, 가져온 데이터가 완전히 복제되었으며 두 클러스터 모두 데이터를 함께 표시할 준비가 되었음을 확인할 수 있습니다.

3단계: 데이터 확인

작업이 ' Completed' 상태에 도달하면, 가져온 엔티티가 두 클러스터 모두에서 표시됩니다. 프라이어리 클러스터에서 컬렉션을 로드하고 쿼리한 다음, 스탠바이 클러스터에서는 컬렉션을 수동으로 로드하지 않은 상태에서 동일한 쿼리를 실행하여 가져온 엔티티가 두 클러스터 모두에 존재하는지 확인하십시오.

스탠바이 클러스터는 스탠바이 상태로 유지되는 동안 읽기 전용입니다. 스탠바이 클러스터에 직접 가져오기, 커밋 또는 기타 DDL/DCL 작업을 제출하지 마십시오. 이러한 작업은 주 클러스터에서 수행하고, CDC 복제를 통해 스탠바이 클러스터에 적용되도록 하십시오.

자주 묻는 질문

가져오기 및 커밋은 어느 클러스터에서 실행해야 합니까?

가져오기 및 커밋은 주 클러스터에서 실행하십시오. 대기 클러스터는 CDC 복제를 통해 가져온 데이터와 커밋을 모두 수신합니다.

스탠바이 클러스터에서 커밋을 수행해야 합니까?

아니요. 프라이머리 클러스터에서 커밋을 수행하면 해당 커밋이 단일 순서 지정 펜스(fence)로 스탠바이 클러스터에 복제됩니다.

왜 ' import in replicating cluster is not supported yet' 오류로 인해 임포트가 실패하나요?

dataCoord.import.enableInReplicatingCluster 해당 클러스터에서 이 기능이 활성화되어 있지 않습니다. 주 클러스터와 대기 클러스터 모두에서 이 기능을 ‘ true ’로 설정하십시오.

auto_commit=true import in replicating cluster is not supported 가 설정된 상태에서 가져오기가 실패하는 이유는 무엇입니까?

복제 클러스터에서는 auto_commit=false 을 사용하는 2PC 가져오기만 허용됩니다. 가져오기 요청에서 options={"auto_commit": "false"} 를 설정하십시오.