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`:
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.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:
| Situasi | Pesan kesalahan |
|---|---|
dataCoord.import.enableInReplicatingCluster belum diaktifkan | import in replicating cluster is not supported yet |
auto_commit=true dikirim | auto_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.