الاستيراد المجمّع في تكرار CDC

يشرح هذا الدليل كيفية تنفيذ عملية استيراد مجمّع على مجموعات Milvus التي تشكل جزءًا من بنية تكرار CDC. في المجموعة التي يتم فيها التكرار، يجب أن تستخدم عملية الاستيراد المجمّع آلية الالتزام ثنائي المراحل (2PC) بحيث يتم الالتزام بعملية الاستيراد كنقطة واحدة مرتبة عبر المجموعة الأساسية والمجموعة الاحتياطية.

في هذا الدليل، تمثل المجموعة الأساسية مجموعة Milvus المصدر، بينما تمثل المجموعة الاحتياطية مجموعة Milvus الهدف.

قبل البدء، تأكد من أن تكرار CDC قد تم تكوينه بالفعل بين مجموعاتك. للحصول على التفاصيل، راجع إعداد تكرار CDC.

لماذا يُعد 2PC ضروريًا

يتم الالتزام التلقائي بعملية الاستيراد المجمّع العادية عند انتهاء مهمة الاستيراد، مما يجعل البيانات المستوردة مرئية على الفور. في بنية تكرار CDC، لا يُسمح بهذا السلوك لأن المجموعتين الرئيسية والاحتياطية يجب أن تجعلا البيانات المستوردة مرئية في نفس النقطة المنطقية.

بدلاً من ذلك، قم بتشغيل الاستيراد في وضع الالتزام ثنائي المراحل عن طريق تعيين « auto_commit=false »:

  1. مرحلة الاستيراد: يقوم Milvus بتحميل البيانات على المجموعة الأساسية ونسخ عملية الاستيراد إلى المجموعة الاحتياطية، لكن البيانات المستوردة تظل غير مرئية. تتوقف مهمة الاستيراد عند حالة "التزام مؤقت" ( Uncommitted ) وتنتظر.

  2. مرحلة الالتزام: تقوم بشكل صريح بالالتزام بمهمة الاستيراد على المجموعة الأساسية. يتم نسخ الالتزام إلى المجموعة الاحتياطية كحاجز واحد مرتب، بحيث تجعل كلتا المجموعتين البيانات المستوردة مرئية في نفس النقطة المنطقية.

الخطوة 1: تمكين الاستيراد في المجموعة المتكررة

يتم تعطيل الاستيراد في المجموعة المتكررة بشكل افتراضي. قم بتمكينه عن طريق تعيين dataCoord.import.enableInReplicatingCluster إلى true في كل من المجموعة الرئيسية والمجموعة الاحتياطية.

إذا قمت بنشر Milvus باستخدام Milvus Operator، فأضف الإعداد التالي إلى spec.config لكل مورد من موارد Milvus:

spec:
  config:
    dataCoord:
      import:
        enableInReplicatingCluster: true

إذا قمت بتكوين Milvus مباشرةً من خلال milvus.yaml ، فأضف الإعداد التالي:

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

قم بتشغيل جميع استدعاءات الاستيراد على المجموعة الأساسية. يتم نسخ البيانات المستوردة وقرار الالتزام تلقائيًا إلى المجموعة الاحتياطية، لذا لا تقم بإرسال أو الالتزام بالاستيراد على المجموعة الاحتياطية بنفسك.

تقوم كل مجموعة بقراءة ملفات الاستيراد من مخزن الكائنات الخاص بها. تأكد من وجود الملفات المراد استيرادها في كل من مخزن الكائنات الأساسي والاحتياطي. يمكنك تحميل الملفات إلى كلتا المجموعتين، أو استخدام مخزن كائنات يمكن لكلتا المجموعتين قراءته. إذا كانت الملفات مفقودة في المجموعة الاحتياطية، فسيفشل الاستيراد المنسوخ هناك مع ظهور خطأ "الكائن غير موجود".

يستخدم المثال التالي أدوات مساعدة الاستيراد المستندة إلى REST من pymilvus.bulk_writer. قيم url هي نفس عناوين Milvus التي تستخدمها لنداءات واجهة برمجة التطبيقات (API) الأخرى.

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.

هل أحتاج إلى إجراء عملية التثبيت على المجموعة الاحتياطية؟

لا. يؤدي إجراء عملية التثبيت على المجموعة الأساسية إلى نسخ عملية التثبيت إلى المجموعة الاحتياطية كحاجز واحد مرتب.

لماذا يفشل الاستيراد مع ظهور الخطأ « import in replicating cluster is not supported yet »؟

dataCoord.import.enableInReplicatingCluster لم يتم تمكين الخيار " true " على تلك المجموعة. اضبطه على " " في كل من المجموعة الرئيسية والمجموعة الاحتياطية.

لماذا يفشل الاستيراد مع " auto_commit=true import in replicating cluster is not supported

في المجموعة التي تتم فيها عملية النسخ المتماثل، لا تُقبل سوى عمليات الاستيراد ذات الخطوتين (2PC) التي تستخدم خيار « auto_commit=false ». قم بتعيين « options={"auto_commit": "false"} » في طلب الاستيراد.