Change Data Capture from MySQL to BigQuery using Debezium and Pub/Sub (Stream) template

תבנית ה-Change Data Capture from MySQL to BigQuery using Debezium and Pub/Sub היא פייפליין סטרימינג שקורא הודעות Pub/Sub עם סימון נתונים שהשתנו (CDC) ממסד נתונים של MySQL וכותב את הרשומות ל-BigQuery. מחבר Debezium מתעד שינויים במסד הנתונים של MySQL ומפרסם את הנתונים ששונו ב-Pub/Sub. לאחר מכן, התבנית קוראת את ההודעות ב-Pub/Sub וכותבת אותן ב-BigQuery.

אפשר להשתמש בתבנית הזו כדי לסנכרן בין מסדי נתונים של MySQL לבין טבלאות של BigQuery. הצינור כותב את הנתונים ששונו לטבלת ביניים ב-BigQuery ומעדכן לסירוגין טבלה ב-BigQuery שמשכפלת את מסד הנתונים של MySQL.

הדרישות לגבי צינורות עיבוד נתונים

  • צריך לפרוס את מחבר Debezium.
  • ההודעות ב-Pub/Sub צריכות להיות מסודרות בסדר מסוים ב-Beam Row.

פרמטרים של תבניות

פרמטרים נדרשים

  • inputSubscriptions: רשימה מופרדת בפסיקים של מינויים לקלט Pub/Sub לקריאה, בפורמט <SUBSCRIPTION_NAME>,<SUBSCRIPTION_NAME>, ....
  • changeLogDataset: מערך הנתונים ב-BigQuery שבו יאוחסנו טבלאות הביניים, בפורמט <DATASET_NAME>.
  • replicaDataset: המיקום של מערך הנתונים ב-BigQuery שבו יאוחסנו טבלאות הרפליקה, בפורמט <DATASET_NAME>.

פרמטרים אופציונליים

  • inputTopics: רשימה מופרדת בפסיקים של נושאי PubSub, שאליהם נדחפים נתוני CDC.
  • updateFrequencySecs: המרווח שבו הפייפליין מעדכן את טבלה ב-BigQuery שמשכפלת את מסד הנתונים של MySQL.
  • useSingleTopic: מגדירים את הערך הזה ל-true אם מגדירים את מחבר Debezium לפרסום כל העדכונים בטבלה בנושא יחיד. ברירת המחדל היא: false.
  • useStorageWriteApi: אם הערך הוא true, צינור הנתונים משתמש ב-BigQuery Storage Write API‏ (https://cloud.google.com/bigquery/docs/write-api). ערך ברירת המחדל הוא false. מידע נוסף זמין במאמר בנושא שימוש ב-Storage Write API‏ (https://beam.apache.org/documentation/io/built-in/google-bigquery/#storage-write-api).
  • useStorageWriteApiAtLeastOnce: כשמשתמשים ב-Storage Write API, המאפיין הזה מציין את סמנטיקת הכתיבה. כדי להשתמש בסמנטיקה של 'לפחות פעם אחת' (https://beam.apache.org/documentation/io/built-in/google-bigquery/#at-least-once-semantics), מגדירים את הפרמטר הזה לערך true. כדי להשתמש בסמנטיקה של מסירה חד-פעמית, מגדירים את הפרמטר לערך false. הפרמטר הזה רלוונטי רק אם הערך של useStorageWriteApi הוא true. ערך ברירת המחדל הוא false.
  • numStorageWriteApiStreams: כשמשתמשים ב-Storage Write API, מציינים את מספר זרמי הכתיבה. אם useStorageWriteApi הוא true ו-useStorageWriteApiAtLeastOnce הוא false, חובה להגדיר את הפרמטר הזה. ברירת המחדל היא 0.
  • storageWriteApiTriggeringFrequencySec: כשמשתמשים ב-Storage Write API, מציינים את תדירות ההפעלה בשניות. אם useStorageWriteApi הוא true ו-useStorageWriteApiAtLeastOnce הוא false, חובה להגדיר את הפרמטר הזה.

הרצת התבנית

כדי להריץ את התבנית הזו:

  1. במחשב המקומי, משכפלים את מאגר DataflowTemplates.
  2. עוברים לספרייה v2/cdc-parent.
  3. מוודאים שחיבור Debezium נפרס.
  4. באמצעות Maven, מריצים את תבנית Dataflow:
    mvn exec:java -pl cdc-change-applier -Dexec.args="--runner=DataflowRunner \
        --inputSubscriptions=SUBSCRIPTIONS \
        --updateFrequencySecs=300 \
        --changeLogDataset=CHANGELOG_DATASET \
        --replicaDataset=REPLICA_DATASET \
        --project=PROJECT_ID \
        --region=REGION_NAME"
      

    מחליפים את מה שכתוב בשדות הבאים:

    • PROJECT_ID: מזהה הפרויקט שבו רוצים להריץ את משימת Dataflow Google Cloud
    • SUBSCRIPTIONS: רשימה מופרדת בפסיקים של שמות מינויים ל-Pub/Sub
    • CHANGELOG_DATASET: מערך הנתונים ב-BigQuery של נתוני יומן השינויים
    • REPLICA_DATASET: מערך הנתונים ב-BigQuery לטבלאות משוכפלות

המאמרים הבאים