תבנית ה-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.
הדרישות לגבי צינורות עיבוד נתונים
פרמטרים של תבניות
פרמטרים נדרשים
- 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, חובה להגדיר את הפרמטר הזה.
הרצת התבנית
כדי להריץ את התבנית הזו:
- במחשב המקומי, משכפלים את מאגר DataflowTemplates.
- עוברים לספרייה
v2/cdc-parent. - מוודאים שחיבור Debezium נפרס.
- באמצעות 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 לטבלאות משוכפלות
-
המאמרים הבאים
- מידע נוסף על תבניות Dataflow
- כאן אפשר לעיין ברשימת התבניות ש-Google סיפקה.