Impor Massal dalam Replikasi CDC

Panduan ini menjelaskan cara menjalankan impor massal pada kluster Milvus yang merupakan bagian dari topologi replikasi CDC. Dalam kluster yang direplikasi, impor massal harus menggunakan dua-fase komit (2PC) agar impor tersebut dikonfirmasi sebagai satu titik yang terurut di seluruh kluster primer dan kluster cadangan.

Dalam panduan ini, kluster primer adalah kluster Milvus sumber, sedangkan kluster siaga adalah kluster Milvus tujuan.

Sebelum memulai, pastikan replikasi CDC telah dikonfigurasi di antara kluster Anda. Untuk detailnya, lihat Mengatur Replikasi CDC.

Mengapa 2PC Diperlukan

Impor massal biasa secara otomatis dikonfirmasi saat pekerjaan impor selesai, yang membuat data yang diimpor langsung terlihat. Dalam topologi replikasi CDC, perilaku ini tidak diperbolehkan karena kluster primer dan siaga harus menampilkan data yang diimpor pada titik logis yang sama.

Sebagai gantinya, jalankan impor dalam mode dua fase komit dengan mengatur ` auto_commit=false`:

  1. Fase impor: Milvus memuat data di kluster primer dan mereplikasi impor ke kluster cadangan, tetapi data yang diimpor tetap tidak terlihat. Tugas impor berhenti pada status " Uncommitted " dan menunggu.

  2. Fase komit: Anda secara eksplisit mengkomit pekerjaan impor pada klaster utama. Komit direplikasi ke klaster siaga sebagai satu pagar terurut, sehingga kedua klaster membuat data yang diimpor terlihat pada titik logis yang sama.

Langkah 1: Aktifkan impor di kluster yang direplikasi

Impor di kluster replikasi dinonaktifkan secara default. Aktifkan dengan mengatur ` dataCoord.import.enableInReplicatingCluster ` menjadi ` true ` pada kluster primer dan kluster siaga.

Jika Anda mengimplementasikan Milvus dengan Milvus Operator, tambahkan pengaturan berikut ke ` spec.config ` pada setiap sumber daya ` Milvus `:

spec:
  config:
    dataCoord:
      import:
        enableInReplicatingCluster: true

Jika Anda mengonfigurasi Milvus secara langsung melalui ` milvus.yaml`, tambahkan pengaturan berikut:

dataCoord:
  import:
    enableInReplicatingCluster: true

Pengaturan ini dapat diperbarui, sehingga dapat berlaku tanpa perlu melakukan restart penuh.

Saat pengaturan ini diaktifkan, kluster replikasi hanya menerima impor dengan ` auto_commit=false`. Tabel berikut mencantumkan permintaan yang umumnya ditolak:

SituasiPesan kesalahan
dataCoord.import.enableInReplicatingCluster belum diaktifkanimport in replicating cluster is not supported yet
auto_commit=true dikirimauto_commit=true import in replicating cluster is not supported

Langkah 2: Jalankan impor 2PC

Jalankan semua panggilan impor pada kluster utama. Data yang diimpor dan keputusan komit akan direplikasi ke kluster cadangan secara otomatis, jadi jangan kirimkan atau lakukan komit impor di kluster cadangan secara manual.

Setiap kluster membaca berkas impor dari penyimpanan objeknya masing-masing. Pastikan berkas yang akan diimpor tersedia di penyimpanan objek primer dan cadangan. Anda dapat mengunggah berkas ke kedua kluster, atau menggunakan penyimpanan objek yang dapat diakses oleh kedua kluster. Jika berkas tidak ada di kluster cadangan, proses impor yang direplikasi akan gagal di sana dengan pesan kesalahan "objek tidak ditemukan".

Contoh berikut menggunakan helper impor berbasis REST dari pymilvus.bulk_writer. Nilai ` url ` adalah alamat Milvus yang sama yang Anda gunakan untuk panggilan API lainnya.

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")

Mengapa harus menunggu hingga Uncommitted di kedua kluster

Melakukan commit sebelum kluster siaga selesai mengimpor tidak akan merusak data, tetapi kluster siaga masih dalam proses mengejar ketertinggalan saat commit diterapkan. Menunggu hingga kluster utama dan siaga sama-sama melaporkan " Uncommitted " memastikan bahwa data yang diimpor telah direplikasi sepenuhnya dan kedua kluster siap menampilkannya secara bersamaan.

Langkah 3: Verifikasi data

Setelah pekerjaan mencapai status " Completed", entitas yang diimpor akan terlihat di kedua kluster. Muat dan jalankan kueri pada koleksi di kluster primer, lalu jalankan kueri yang sama di kluster cadangan tanpa memuat koleksi tersebut secara manual di sana, dan pastikan entitas yang diimpor terdapat di kedua kluster.

Cluster siaga bersifat read-only selama masih berstatus siaga. Jangan mengirimkan impor, commit, atau operasi DDL atau DCL lainnya secara langsung ke cluster siaga. Lakukan operasi ini di cluster utama dan biarkan replikasi CDC menerapkannya ke cluster siaga.

FAQ

Di kluster mana saya harus menjalankan impor dan commit?

Jalankan impor dan komit di klaster utama. Klaster siaga menerima data yang diimpor dan komit melalui replikasi CDC.

Apakah saya perlu melakukan commit di kluster standby?

Tidak. Melakukan commit di kluster primer akan mereplikasi commit tersebut ke kluster siaga sebagai satu fence yang terurut.

Mengapa proses impor saya gagal dengan pesan " import in replicating cluster is not supported yet"?

dataCoord.import.enableInReplicatingCluster tidak diaktifkan pada cluster tersebut. Atur menjadi " true " pada cluster utama dan cluster cadangan.

Mengapa impor saya gagal dengan pesan " auto_commit=true import in replicating cluster is not supported"?

Dalam kluster replikasi, hanya impor 2PC dengan opsi ` auto_commit=false ` yang diterima. Atur ` options={"auto_commit": "false"} ` pada permintaan impor.