תבנית Apache Kafka ל-BigQuery היא פייפליין להזרמת נתונים שקולט נתוני טקסט מאשכולות של שירות מנוהל של Google Cloud ל-Apache Kafka, ואז מוציא את הרשומות שמתקבלות לטבלאות ב-BigQuery. כל השגיאות שמתרחשות במהלך הוספת הנתונים לטבלת הפלט מתווספות לטבלת שגיאות נפרדת ב-BigQuery.
אפשר גם להשתמש בתבנית Apache Kafka to BigQuery עם Kafka בניהול עצמי או חיצוני.
הדרישות לגבי צינורות עיבוד נתונים
- שרת הברוקר של Apache Kafka צריך לפעול ולהיות נגיש ממכונות העובד של Dataflow.
- נושאי Apache Kafka חייבים להתקיים.
- צריך להפעיל את ממשקי ה-API של Dataflow, BigQuery ו-Cloud Storage. אם נדרש אימות, צריך גם להפעיל את Secret Manager API.
- יוצרים מערך נתונים וטבלה ב-BigQuery עם הסכימה המתאימה לנושא הקלט של Kafka. אם משתמשים בכמה סכימות באותו נושא ורוצים לכתוב לכמה טבלאות, לא צריך ליצור את הטבלה לפני שמגדירים את צינור הנתונים.
- כשמפעילים את תור ההודעות שלא עברו עיבוד (dead-letter) של התבנית, צריך ליצור טבלה ריקה בלי סכימה לתור ההודעות שלא עברו עיבוד.
- אם מתחברים לאשכול של שירות מנוהל ל-Apache Kafka, הצינור צריך לעמוד גם בדרישות שמפורטות במאמר שימוש ב-Dataflow עם שירות מנוהל ל-Apache Kafka.
פורמט ההודעה ב-Kafka
התבנית הזו תומכת בקריאת הודעות מ-Kafka בפורמטים הבאים:
פורמט JSON
כדי לקרוא הודעות JSON, מגדירים את פרמטר התבנית messageFormat לערך "JSON".
קידוד בינארי של Avro
כדי לקרוא הודעות בינאריות בפורמט Avro, מגדירים את פרמטרים התבנית הבאים:
messageFormat:"AVRO_BINARY_ENCODING".-
binaryAvroSchemaPath: המיקום של קובץ סכימת Avro ב-Cloud Storage. דוגמה:gs://BUCKET_NAME/message-schema.avsc.
מידע נוסף על פורמט בינארי של Avro זמין במאמר קידוד בינארי במאמרי העזרה של Apache Avro.
Avro עם קידוד של Confluent Schema Registry
כדי לקרוא הודעות ב-Avro שמקודדות ב-Confluent Schema Registry, מגדירים את פרמטרים התבנית הבאים:
messageFormat:"AVRO_CONFLUENT_WIRE_FORMAT".-
schemaFormat: אחד מהערכים הבאים:-
"SINGLE_SCHEMA_FILE": סכימת ההודעה מוגדרת בקובץ סכימת Avro. מציינים את המיקום של קובץ הסכימה ב-Cloud Storage בפרמטרconfluentAvroSchemaPath. -
"SCHEMA_REGISTRY": ההודעות מקודדות באמצעות Confluent Schema Registry. מציינים את כתובת ה-URL של מופע Confluent Schema Registry בפרמטרschemaRegistryConnectionUrl, ומציינים את מצב האימות בפרמטרschemaRegistryAuthenticationMode.
-
מידע נוסף על הפורמט הזה זמין במאמר בנושא פורמט העברה במאמרי העזרה של Confluent.
אימות
תבנית Apache Kafka ל-BigQuery תומכת באימות SASL/PLAIN לשרתי Kafka.
פרמטרים של תבניות
פרמטרים נדרשים
- readBootstrapServerAndTopic: נושא Kafka לקריאת הקלט ממנו.
- writeMode: כתיבת רשומות לטבלה אחת או לכמה טבלאות (על סמך הסכימה). מצב
DYNAMIC_TABLE_NAMESנתמך רק בAVRO_CONFLUENT_WIRE_FORMATפורמט הודעת מקור ובSCHEMA_REGISTRYמקור סכימה. שם טבלת היעד נוצר באופן אוטומטי על סמך שם סכימת Avro של כל הודעה. יכולה להיות סכימה אחת (שיוצרת טבלה אחת) או כמה סכימות (שיוצרות כמה טבלאות). במצבSINGLE_TABLE_NAME, המערכת כותבת לטבלה אחת (סכימה אחת) שהמשתמש מציין. ברירת המחדל היאSINGLE_TABLE_NAME. - kafkaReadAuthenticationMode: מצב האימות לשימוש עם אשכול Kafka. משתמשים ב-
NONEלאימות, ב-SASL_PLAINלשם משתמש וסיסמה של SASL/PLAIN, ב-SASL_SCRAM_512לאימות SASL_SCRAM_512 וב-TLSלאימות מבוסס-אישור. המאפייןAPPLICATION_DEFAULT_CREDENTIALSצריך לשמש רק עבור אשכול של שירות מנוהל של Google Cloud ל-Apache Kafka. הוא מאפשר אימות באמצעות פרטי כניסה שמוגדרים כברירת מחדל באפליקציה. - messageFormat: הפורמט של הודעות Kafka לקריאה. הערכים הנתמכים הם
AVRO_CONFLUENT_WIRE_FORMAT(Avro מקודד של Confluent Schema Registry),AVRO_BINARY_ENCODING(Avro בינארי רגיל) ו-JSON. ברירת המחדל היא: AVRO_CONFLUENT_WIRE_FORMAT. - useBigQueryDLQ: אם הערך הוא true, הודעות שנכשלו ייכתבו ל-BigQuery עם מידע נוסף על השגיאה. ברירת המחדל היא: false.
פרמטרים אופציונליים
- outputTableSpec: מיקום הטבלה ב-BigQuery שאליה ייכתב הפלט. השם צריך להיות בפורמט
<project>:<dataset>.<table_name>. הסכימה של הטבלה צריכה להתאים לאובייקטים של הקלט. - persistKafkaKey: אם הערך הוא true, הצינור ישמור את מפתח ההודעה של Kafka בטבלה ב-BigQuery, בשדה
_keyמסוגBYTES. ברירת המחדל היאfalse(המערכת מתעלמת מהמפתח). - outputProject: פרויקט הפלט ב-BigQuery שבו נמצא מערך הנתונים. הטבלאות ייווצרו באופן דינמי במערך הנתונים. ברירת המחדל היא ריק.
- outputDataset: מערך נתוני הפלט של BigQuery שאליו ייכתב הפלט. הטבלאות ייווצרו באופן דינמי במערך הנתונים. אם הטבלאות נוצרות מראש, השמות שלהן צריכים להיות בהתאם למוסכמות למתן שמות שצוינו. השם צריך להיות
bqTableNamePrefix + Avro Schema FullName, כל מילה תופרד באמצעות מקף-. ברירת המחדל היא ריק. - bqTableNamePrefix: קידומת לשם שתשמש ליצירת טבלאות פלט ב-BigQuery. ההגדרה הזו רלוונטית רק כשמשתמשים במאגר סכימות. ברירת המחדל היא ריק.
- createDisposition: BigQuery CreateDisposition. לדוגמה:
CREATE_IF_NEEDED, CREATE_NEVER. ברירת המחדל היא: CREATE_IF_NEEDED. - writeDisposition: BigQuery WriteDisposition. לדוגמה:
WRITE_APPEND, WRITE_EMPTYאוWRITE_TRUNCATE. ברירת המחדל היא: WRITE_APPEND. - useAutoSharding: אם הערך הוא true, הצינור משתמש בפיצול אוטומטי של נתונים (auto-sharding) כשמתבצעת כתיבה ל-BigQuery. ערך ברירת המחדל הוא
true. - numStorageWriteApiStreams: מציין את מספר זרמי הכתיבה. חובה להגדיר את הפרמטר הזה. ערך ברירת המחדל הוא
0. - storageWriteApiTriggeringFrequencySec: הפרמטר הזה מציין את תדירות ההפעלה בשניות, וחובה להגדיר אותו. ברירת המחדל היא 5 שניות.
- useStorageWriteApiAtLeastOnce: הפרמטר הזה תקף רק אם האפשרות 'שימוש ב-BigQuery Storage Write API' מופעלת. אם האפשרות הזו מופעלת, נעשה שימוש בסמנטיקה של מסירה לפחות פעם אחת עבור Storage Write API. אחרת, נעשה שימוש בסמנטיקה של מסירה חד-פעמית. ברירת המחדל היא: false.
- enableCommitOffsets: ביצוע Commit של אופסטים של הודעות שעברו עיבוד ל-Kafka. אם ההגדרה הזו מופעלת, היא תמזער את הפערים או את העיבוד הכפול של הודעות כשמפעילים מחדש את הצינור. צריך לציין את המזהה של קבוצת הצרכנים. ברירת המחדל היא: false.
- consumerGroupId: המזהה הייחודי של קבוצת הצרכנים שאליה הצינור הזה משתייך. חובה אם האפשרות Commit Offsets to Kafka (שמירת היסטוריית המיקומים ב-Kafka) מופעלת. ברירת המחדל היא ריק.
- kafkaReadOffset: נקודת ההתחלה לקריאת הודעות כשאין היסטים מחויבים. ההתחלה המוקדמת ביותר היא מההתחלה, וההתחלה המאוחרת ביותר היא מההודעה החדשה ביותר. ברירת המחדל היא: latest.
- kafkaReadUsernameSecretId: מזהה הסוד ב-Google Cloud Secret Manager שמכיל את שם המשתמש ב-Kafka שמשמש לאימות
SASL_PLAIN. לדוגמה,projects/<PROJECT_ID>/secrets/<SECRET_ID>/versions/<SECRET_VERSION>. ברירת המחדל היא ריק. - kafkaReadPasswordSecretId: מזהה הסוד ב-Google Cloud Secret Manager שמכיל את הסיסמה ל-Kafka שבה יש להשתמש עם אימות
SASL_PLAIN. לדוגמה,projects/<PROJECT_ID>/secrets/<SECRET_ID>/versions/<SECRET_VERSION>. ברירת המחדל היא ריק. - kafkaReadKeystoreLocation: הנתיב ב-Google Cloud Storage לקובץ Java KeyStore (JKS) שמכיל את אישור ה-TLS והמפתח הפרטי שבהם יש להשתמש בעת אימות מול אשכול Kafka. לדוגמה,
gs://your-bucket/keystore.jks. - kafkaReadTruststoreLocation: הנתיב ב-Google Cloud Storage לקובץ Java TrustStore (JKS) שמכיל את האישורים המהימנים שמשמשים לאימות הזהות של ברוקר Kafka.
- kafkaReadTruststorePasswordSecretId: מזהה הסוד ב-Google Cloud Secret Manager שמכיל את הסיסמה לגישה לקובץ Java TrustStore (JKS) לאימות TLS של Kafka. לדוגמה,
projects/<PROJECT_ID>/secrets/<SECRET_ID>/versions/<SECRET_VERSION>. - kafkaReadKeystorePasswordSecretId: מזהה הסוד ב-Google Cloud Secret Manager שמכיל את הסיסמה לגישה לקובץ Java KeyStore (JKS) לאימות Kafka TLS. לדוגמה,
projects/<PROJECT_ID>/secrets/<SECRET_ID>/versions/<SECRET_VERSION>. - kafkaReadKeyPasswordSecretId: מזהה הסוד ב-Google Cloud Secret Manager שמכיל את הסיסמה לגישה למפתח הפרטי בקובץ Java KeyStore (JKS) לאימות Kafka TLS. לדוגמה,
projects/<PROJECT_ID>/secrets/<SECRET_ID>/versions/<SECRET_VERSION>. - kafkaReadSaslScramUsernameSecretId: המזהה הסודי של Google Cloud Secret Manager שמכיל את שם המשתמש של Kafka לשימוש באימות
SASL_SCRAM. לדוגמה,projects/<PROJECT_ID>/secrets/<SECRET_ID>/versions/<SECRET_VERSION>. - kafkaReadSaslScramPasswordSecretId: מזהה הסוד ב-Google Cloud Secret Manager שמכיל את הסיסמה ל-Kafka שבה יש להשתמש עם אימות
SASL_SCRAM. לדוגמה,projects/<PROJECT_ID>/secrets/<SECRET_ID>/versions/<SECRET_VERSION>. - kafkaReadSaslScramTruststoreLocation: הנתיב ב-Google Cloud Storage לקובץ Java TrustStore (JKS) שמכיל את האישורים המהימנים שמשמשים לאימות הזהות של ברוקר Kafka.
- kafkaReadSaslScramTruststorePasswordSecretId: מזהה הסוד ב-Google Cloud Secret Manager שמכיל את הסיסמה לגישה לקובץ Java TrustStore (JKS) לאימות Kafka SASL_SCRAM. לדוגמה,
projects/<PROJECT_ID>/secrets/<SECRET_ID>/versions/<SECRET_VERSION>. - schemaFormat: פורמט הסכימה של Kafka. אפשר לציין את הערך כ-
SINGLE_SCHEMA_FILEאו כ-SCHEMA_REGISTRY. אם מצויןSINGLE_SCHEMA_FILE, המערכת משתמשת בסכימה שמוזכרת בקובץ הסכימה של avro לכל ההודעות. אם מציינים אתSCHEMA_REGISTRY, ההודעות יכולות להיות עם סכימה אחת או עם כמה סכימות. ברירת המחדל היא: SINGLE_SCHEMA_FILE. - confluentAvroSchemaPath: הנתיב ב-Google Cloud Storage לקובץ סכימת Avro יחיד שמשמש לפענוח כל ההודעות בנושא. ברירת המחדל היא ריק.
- schemaRegistryConnectionUrl: כתובת ה-URL של מופע Confluent Schema Registry שמשמש לניהול סכימות Avro לפענוח הודעות. ברירת המחדל היא ריק.
- binaryAvroSchemaPath: הנתיב ב-Google Cloud Storage לקובץ סכימת Avro שמשמש לפענוח הודעות Avro עם קידוד בינארי. ברירת המחדל היא ריק.
- schemaRegistryAuthenticationMode: מצב האימות לשימוש עם Schema Registry. משתמשים ב-
NONEללא אימות, ב-TLSלאימות מבוסס-אישור וב-OAUTHלאימות OAuth2. APPLICATION_DEFAULT_CREDENTIALSצריך לשמש רק לשירות מנוהל של Google Cloud ל-Apache Kafka Schema Registry, והוא מאפשר אימות באמצעות Application Default Credentials. - schemaRegistryTruststoreLocation: המיקום של אישור ה-SSL שבו מאוחסן מאגר האישורים לאימות ב-Schema Registry. לדוגמה,
/your-bucket/truststore.jks. - schemaRegistryTruststorePasswordSecretId: SecretId במנהל הסודות שבו מאוחסנת הסיסמה לגישה לסוד ב-truststore. לדוגמה,
projects/your-project-number/secrets/your-secret-name/versions/your-secret-version. - schemaRegistryKeystoreLocation: מיקום Keystore שמכיל את אישור ה-SSL והמפתח הפרטי. לדוגמה,
/your-bucket/keystore.jks. - schemaRegistryKeystorePasswordSecretId: SecretId ב-Secret Manager שבו הסיסמה לגישה לקובץ של מאגר המפתחות. לדוגמה,
projects/your-project-number/secrets/your-secret-name/versions/your-secret-version. - schemaRegistryKeyPasswordSecretId: SecretId של הסיסמה שנדרשת לגישה למפתח הפרטי של הלקוח שמאוחסן במאגר המפתחות. לדוגמה,
projects/your-project-number/secrets/your-secret-name/versions/your-secret-version. - schemaRegistryOauthClientId: מזהה הלקוח שמשמש לאימות לקוח Schema Registry במצב OAUTH. נדרש לפורמט ההודעה AVRO_CONFLUENT_WIRE_FORMAT.
- schemaRegistryOauthClientSecretId: מזהה הסוד ב-Google Cloud Secret Manager שמכיל את סוד הלקוח שמשמש לאימות לקוח Schema Registry במצב OAUTH. נדרש לפורמט ההודעה AVRO_CONFLUENT_WIRE_FORMAT. לדוגמה,
projects/<PROJECT_ID>/secrets/<SECRET_ID>/versions/<SECRET_VERSION>. - schemaRegistryOauthScope: היקף אסימון הגישה שמשמש לאימות לקוח Schema Registry במצב OAUTH. השדה הזה הוא אופציונלי, כי אפשר לשלוח את הבקשה בלי להעביר פרמטר היקף. לדוגמה,
openid. - schemaRegistryOauthTokenEndpointUrl: כתובת ה-URL שמבוססת על HTTP(S) של ספק הזהויות OAuth/OIDC שמשמש לאימות לקוח Schema Registry במצב OAUTH. נדרש לפורמט ההודעה AVRO_CONFLUENT_WIRE_FORMAT.
- outputDeadletterTable: שם הטבלה ב-BigQuery שכולל את שם הפרויקט, מערך הנתונים והטבלה, שאליה מועברות הודעות שנכשלו. הודעות שלא הגיעו לטבלת הפלט מסיבות שונות (למשל, סכימה לא תואמת, JSON לא תקין) נכתבות בטבלה הזו. הטבלה נוצרת על ידי התבנית. לדוגמה,
your-project-id:your-dataset.your-table-name. - javascriptTextTransformGcsPath: ה-URI של Cloud Storage של קובץ ה- .js שמגדיר את הפונקציה בהגדרת המשתמש (UDF) ב-JavaScript שבה רוצים להשתמש. לדוגמה,
gs://my-bucket/my-udfs/my_file.js. - javascriptTextTransformFunctionName: השם של פונקציית JavaScript בהגדרת המשתמש (UDF) שבה רוצים להשתמש. לדוגמה, אם קוד פונקציית ה-JavaScript הוא
myTransform(inJson) { /*...do stuff...*/ }, אז שם הפונקציה הואmyTransform. דוגמאות ל-UDF ב-JavaScript זמינות במאמר UDF Examples (https://github.com/GoogleCloudPlatform/DataflowTemplates#udf-examples). - javascriptTextTransformReloadIntervalMinutes: מציין באיזו תדירות לטעון מחדש את ה-UDF, בדקות. אם הערך גדול מ-0, מערכת Dataflow בודקת מעת לעת את קובץ ה-UDF ב-Cloud Storage, ומטעינה מחדש את ה-UDF אם הקובץ משתנה. הפרמטר הזה מאפשר לכם לעדכן את הפונקציה המוגדרת על ידי המשתמש בזמן שהצינור פועל, בלי שתצטרכו להפעיל מחדש את העבודה. אם הערך הוא
0, טעינה מחדש של פונקציות UDF מושבתת. ערך ברירת המחדל הוא0.
פונקציה בהגדרת המשתמש
אפשר גם להרחיב את התבנית הזו על ידי כתיבת פונקציה בהגדרת המשתמש (UDF). התבנית קוראת ל-UDF לכל רכיב קלט. מטענים ייעודיים של רכיבים עוברים סריאליזציה כמחרוזות JSON. למידע נוסף, ראו יצירת פונקציות מוגדרות על ידי המשתמש לתבניות Dataflow.
התבנית תומכת בפונקציות UDF רק בהודעות Kafka בפורמט JSON. אם ההודעות ב-Kafka הן בפורמט Avro, הפונקציה המוגדרת על ידי המשתמש לא מופעלת.מפרט הפונקציה
המאפיינים של פונקציית UDF:
- קלט: ערך הרשומה ב-Kafka, שעבר סריאליזציה כמחרוזת JSON
- פלט: מחרוזת JSON שתואמת לסכימה של טבלת היעד ב-BigQuery
הרצת התבנית
המסוף
- עוברים לדף Create job from template (יצירת משימה מתבנית) ב-Dataflow. כניסה לדף Create job from template
- בשדה שם המשימה, מזינים שם ייחודי למשימה.
- אופציונלי: בשדה Regional endpoint (נקודת קצה אזורית), בוחרים ערך מהתפריט הנפתח. אזור ברירת המחדל הוא
us-central1.רשימת האזורים שבהם אפשר להריץ משימת Dataflow מופיעה במאמר בנושא מיקומי Dataflow.
- בתפריט הנפתח Dataflow template (תבנית Dataflow), בוחרים את התבנית Kafka to BigQuery (מ-Kafka ל-BigQuery).
- בשדות הפרמטרים שמופיעים, מזינים את ערכי הפרמטרים.
- אופציונלי: כדי לעבור מעיבוד של כל נתון בדיוק פעם אחת אל מצב סטרימינג של כל נתון לפחות פעם אחת, בוחרים באפשרות לפחות פעם אחת.
- לוחצים על הפעלת העבודה.
gcloud
במעטפת או בטרמינל, מריצים את התבנית:
gcloud dataflow flex-template run JOB_NAME \ --project=PROJECT_ID \ --region=REGION_NAME \ --template-file-gcs-location=gs://dataflow-templates-REGION_NAME/VERSION/flex/ \ --parameters \ readBootstrapServerAndTopic=BOOTSTRAP_SERVER_AND_TOPIC,\ kafkaReadAuthenticationMode=APPLICATION_DEFAULT_CREDENTIALS,\ messageFormat=JSON,\ writeMode=SINGLE_TABLE_NAME,\ outputTableSpec=PROJECT_ID:DATASET_NAME.TABLE_NAME,\ useBigQueryDLQ=true,\ outputDeadletterTable=PROJECT_ID:DATASET_NAME.ERROR_TABLE_NAME
מחליפים את מה שכתוב בשדות הבאים:
-
PROJECT_ID: מזהה הפרויקט שבו רוצים להריץ את משימת Dataflow Google Cloud -
JOB_NAME: שם ייחודי של המשימה לפי בחירתכם -
REGION_NAME: האזור שבו רוצים לפרוס את עבודת Dataflow, לדוגמה:us-central1 -
VERSION: הגרסה של התבנית שבה רוצים להשתמשאפשר להשתמש בערכים הבאים:
-
latestכדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית האב ללא תאריך בדלי – gs://dataflow-templates-REGION_NAME/latest/ - שם הגרסה, כמו
2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית, שאפשר למצוא אותה בתיקיית האב המתאימה עם התאריך בדלי – gs://dataflow-templates-REGION_NAME/
-
-
BOOTSTRAP_SERVER_AND_TOPIC: כתובת שרת ה-bootstrap והנושא של Apache Kafkaהפורמט של כתובת שרת האתחול והנושא תלוי בסוג האשכול:
- קלאסטר שירות מנוהל ל-Apache Kafka:
projects/PROJECT_ID/locations/REGION_NAME/clusters/CLUSTER_NAME/topics/TOPIC_NAME - אשכול Kafka חיצוני:
BOOTSTRAP_SERVER_ADDRESS;TOPIC_NAME
- קלאסטר שירות מנוהל ל-Apache Kafka:
-
DATASET_NAME: השם של מערך הנתונים ב-BigQuery -
TABLE_NAME: השם של טבלת הפלט ב-BigQuery -
ERROR_TABLE_NAME: השם של הטבלה ב-BigQuery שאליה ייכתבו רשומות השגיאות
API
כדי להריץ את התבנית באמצעות API בארכיטקטורת REST, שולחים בקשת HTTP POST. מידע נוסף על ה-API ועל היקפי ההרשאות שלו זמין במאמר projects.templates.launch.
POST https://dataflow.googleapis.com/v1b3/projects/PROJECT_ID/locations/LOCATION/flexTemplates:launch { "launch_parameter": { "jobName": "JOB_NAME", "parameters": { "readBootstrapServerAndTopic": "BOOTSTRAP_SERVER_AND_TOPIC", "kafkaReadAuthenticationMode": "APPLICATION_DEFAULT_CREDENTIALS", "messageFormat": "JSON", "writeMode": "SINGLE_TABLE_NAME", "outputTableSpec": "PROJECT_ID:DATASET_NAME.TABLE_NAME", "useBigQueryDLQ": "true", "outputDeadletterTable": "PROJECT_ID:DATASET_NAME.ERROR_TABLE_NAME" }, "containerSpecGcsPath": "gs://dataflow-templates-LOCATION/VERSION/flex/", } }
מחליפים את מה שכתוב בשדות הבאים:
-
PROJECT_ID: מזהה הפרויקט שבו רוצים להריץ את משימת Dataflow Google Cloud -
JOB_NAME: שם ייחודי של המשימה לפי בחירתכם -
LOCATION: האזור שבו רוצים לפרוס את עבודת Dataflow, לדוגמה:us-central1 -
VERSION: הגרסה של התבנית שבה רוצים להשתמשאפשר להשתמש בערכים הבאים:
-
latestכדי להשתמש בגרסה העדכנית של התבנית, שזמינה בתיקיית האב ללא תאריך בדלי – gs://dataflow-templates-REGION_NAME/latest/ - שם הגרסה, כמו
2023-09-12-00_RC00, כדי להשתמש בגרסה ספציפית של התבנית, שאפשר למצוא אותה בתיקיית האב המתאימה עם התאריך בדלי – gs://dataflow-templates-REGION_NAME/
-
-
BOOTSTRAP_SERVER_AND_TOPIC: כתובת שרת ה-bootstrap והנושא של Apache Kafkaהפורמט של כתובת שרת האתחול והנושא תלוי בסוג האשכול:
- קלאסטר שירות מנוהל ל-Apache Kafka:
projects/PROJECT_ID/locations/LOCATION/clusters/CLUSTER_NAME/topics/TOPIC_NAME - אשכול Kafka חיצוני:
BOOTSTRAP_SERVER_ADDRESS;TOPIC_NAME
- קלאסטר שירות מנוהל ל-Apache Kafka:
-
DATASET_NAME: השם של מערך הנתונים ב-BigQuery -
TABLE_NAME: השם של טבלת הפלט ב-BigQuery -
ERROR_TABLE_NAME: השם של הטבלה ב-BigQuery שאליה ייכתבו רשומות השגיאות
מידע נוסף זמין במאמר כתיבת נתונים מ-Kafka ל-BigQuery באמצעות Dataflow.
המאמרים הבאים
- מידע נוסף על תבניות Dataflow
- כאן אפשר לעיין ברשימת התבניות ש-Google סיפקה.