אם אתם צריכים לטעון נתונים מנושא של שירות מנוהל של Google Cloud ל-Apache Kafka לקטגוריה של Cloud Storage, אתם יכולים לעשות זאת באמצעות תבנית Dataflow. אפשר להשתמש במסוף Google Cloud , ב-API בארכיטקטורת REST או ב-Google Cloud CLI.
במאמר הזה מוסבר איך להגדיר את תבנית Kafka to Cloud Storage Dataflow באמצעות מסוף Google Cloud .
Google Cloud מוצרים בשימוש
בתבנית Kafka to Cloud Storage Dataflow נעשה שימוש במוצרים הבאים שחלים עליהם חיובים Google Cloud . אתם יכולים להשתמש במחשבון התמחור כדי ליצור הערכת עלויות בהתאם לשימוש החזוי.
- Dataflow: Dataflow הוא שירות מנוהל במלואו לעיבוד נתונים. התבנית Kafka to Cloud Storage Dataflow משתמשת ב-Dataflow כדי ליצור צינור שקורא נתונים מנושא Kafka, מבצע את כל השינויים הנדרשים וכותב אותם ל-Cloud Storage. התכונות של Dataflow להתאמה אוטומטית לעומס ולתיקון עצמי מבטיחות שצינור הנתונים יפעל בצורה אמינה ויעילה.
- Cloud Storage: משמש כיעד לנתוני Kafka. תצטרכו קטגוריה של Cloud Storage כדי לאחסן את הנתונים שמועברים על ידי צינור הנתונים של Dataflow.
בנוסף, הפתרון משתמש גם בשירות המנוהל של Google Cloud ל-Apache Kafka.
- שירות מנוהל של Google Cloud ל-Apache Kafka: שירות Google Cloud שעוזר להפעיל את Apache Kafka. מכיל את נתוני המקור של הפייפליין. תצטרכו נושא ואשכול קיימים של השירות המנוהל ל-Apache Kafka עם נתונים שאתם רוצים להעביר ל-Cloud Storage. מידע נוסף על התמחור של שירות מנוהל של Google Cloud ל-Apache Kafka זמין במדריך התמחור.
לפני שמתחילים
לפני שמפעילים את תבנית Kafka to Cloud Storage Dataflow, צריך לוודא שביצעתם את הפעולות הבאות:
יוצרים אשכול ונושא בשירות מנוהל ל-Apache Kafka.
אחת הדרכים ליצור אשכול ונושא היא לפעול לפי ההוראות שבמדריך למתחילים של שירות מנוהל ל-Apache Kafka.
אם הנושא מכיל רשומות Avro, אפשר לעיין במאמר ציון פורמט ההודעה כדי לראות מהן דרישות המשאבים הנוספות.
מפעילים את ממשקי ה-API הבאים: Google Cloud
Dataflow
Cloud Storage
gcloud services enable dataflow.googleapis.com storage-api.googleapis.com \יצירת קטגוריה של Cloud Storage.
מידע נוסף על יצירת קטגוריה של Cloud Storage זמין במאמר יצירת קטגוריה.
מקצים את התפקיד Managed Kafka client לחשבון השירות של העובד (worker) ב-Dataflow
כדי לקשר את משימת Dataflow לשירות המנוהל ל-Apache Kafka, צריך להעניק הרשאות ספציפיות לחשבון השירות של עובד Dataflow. חשבון השירות הזה הוא הזהות שמשמשת את כל מכונות ה-VM של העובדים בעבודת Dataflow, וכל הבקשות שנשלחות ממכונות ה-VM האלה משתמשות בחשבון הזה.
כדי לאפשר גישה למשאבי Kafka, צריך להעניק את התפקיד roles/managedkafka.client לחשבון השירות של Dataflow worker. התפקיד הזה כולל את ההרשאה managedkafka.clusters.connect שנדרשת ליצירת חיבורים.
מידע נוסף על חשבון השירות של העובד זמין במאמר אבטחה והרשאות של צינורות ב-Google Cloud.
כדי להעניק את התפקיד Managed Kafka client לחשבון השירות של Dataflow:
המסוף
- נכנסים לדף IAM במסוף Google Cloud .
כניסה לדף IAM - מוודאים שהפרויקט מוגדר כפרויקט הצרכן שאליו יתבצעו גישות של לקוח השירות המנוהל ל-Apache Kafka.
- לוחצים על הענקת גישה.
- בדף החדש, בשדה Add Principals, מזינים את כתובת האימייל של חשבון השירות של עובד Dataflow שבו אתם משתמשים.
- בקטע Assign roles (הקצאת תפקידים), בוחרים בתפקיד Managed Kafka client (לקוח מנוהל של Kafka).
- לוחצים על Save.
CLI של gcloud
-
במסוף Google Cloud , מפעילים את Cloud Shell.
בחלק התחתון של Google Cloud המסוף יתחיל סשן של Cloud Shell ותופיע הודעה של שורת הפקודה. Cloud Shell היא סביבת מעטפת שבה ה-CLI של Google Cloud מותקן ומוגדרים ערכים לפרויקט הקיים. הסשן יופעל תוך כמה שניות.
-
מריצים את הפקודה
gcloud projects add-iam-policy-binding:gcloud projects add-iam-policy-binding PROJECT_ID \ --member serviceAccount:SERVICE_ACCOUNT_EMAIL \ --role roles/managedkafka.client
מחליפים את מה שכתוב בשדות הבאים:
-
PROJECT_ID הוא מזהה הפרויקט.
-
SERVICE_ACCOUNT_EMAIL היא כתובת האימייל של חשבון השירות של העובד (worker) ב-Dataflow.
-
הפעלת תבנית Kafka ל-Cloud Storage Dataflow
אפשר להפעיל את תבנית Kafka to Cloud Storage Dataflow מדף פרטי האשכול במסוף.
-
נכנסים לדף Cluster במסוף Google Cloud .
מוצגת רשימה של האשכולות שיצרתם בפרויקט.
- כדי לראות את דף הפרטים של האשכול, לוחצים על שם האשכול.
- בדף הפרטים של האשכול, לוחצים על ייבוא נתונים.
ייפתח הדף Create a Dataflow job using template "Kafka to Kafka" (יצירת משימת Dataflow באמצעות התבנית Kafka to Kafka).
- בתבנית, בשדה Dataflow template, מעדכנים את התבנית ל-Kafka to Cloud Storage.
מגדירים את השדות בתבנית לפי המידע שמופיע בקטעים הבאים.
מזינים שם של משרה
בשדה Job name (שם המשימה), מזינים שם למשימת Dataflow.
השם חייב להיות ייחודי בין כל המשימות שפועלות כרגע בפרויקט.
בחירת נקודת קצה אזורית לצינור
בשדה Regional endpoint, מגדירים את נקודת הקצה האזורית למיקום של אשכול Kafka כדי לצמצם את העמלות על העברת נתונים בין אזורים.
העובדים של Dataflow יכולים לפעול באופן עצמאי מאזור האשכול של Kafka. עם זאת, אם מפעילים עובדים מחוץ לאזור של אשכול Kafka, נושאים בעלויות של תעבורת נתונים יוצאת בין אזורים.
כדי לראות את המיקום של האשכול, פועלים לפי השלבים במאמר בנושא הצגת רשימה של אשכולות בשירות המנוהל ל-Apache Kafka.
הגדרת המקור
בקטע מקור, משאירים את ערך ברירת המחדל שירות מנוהל ל-Apache Kafka.
בקטעים Kafka cluster ו-Kafka source authentication mode, משאירים את ערכי ברירת המחדל.
בקטע Kafka topic (נושא Kafka), בוחרים נושא מתוך רשימת הנושאים הזמינים.
הגדרת פורמט ההודעה של Kafka
תבנית Dataflow תומכת בשלושת הפורמטים הבאים של הודעות:
פורמט Avro Confluent wire: כל הודעת Kafka כוללת בייט magic, מזהה סכימה ורשומה בקידוד בינארי של Avro.
בפורמטים של Avro (פורמט Confluent wire), אפשר להשתמש בסכימה אחת או בכמה סכימות:
סכימה אחת: כל ההודעות תואמות לסכימת Avro מוגדרת מראש.
סכימות מרובות: הודעות יכולות להשתמש בסכימות שונות. התמיכה הזו זמינה רק בפורמט Avro (פורמט Confluent wire).
Avro (קידוד בינארי): ההודעות מכילות רק את המטען הייעודי (payload) של הרשומה, ללא מטא-נתונים. צריך לספק קובץ סכמת Avro (.avsc) שהועלה ל-Cloud Storage. כל ההודעות צריכות להיות בהתאם לסכימה היחידה הזו.
JSON: לא נדרשת סכימה מוגדרת לרשומות. רשומות שלא תואמות לסכימה נשלחות לתור של הודעות שלא ניתן למסור (אם הוא מוגדר) או שמתבצע רישום של הודעת שגיאה. הפורמט הנתמך הוא פורמט
{"field": "value"}. הפורמט[{"name": "field", "value": "value"}]לא נתמך.
שירות מנוהל של Google Cloud ל-Apache Kafka לא מציע מאגר סכימות. התבנית תומכת רק בהעברת פרטי אימות למאגרי סכימות שתואמים לפורמט Confluent-wire.
פורמט Avro Confluent wire
אם בוחרים באפשרות הזו כפורמט ההודעה של Kafka, צריך להגדיר את ההגדרות הנוספות הבאות:
מקור הסכימה: בשדה הזה מציינים את המקום שבו נמצאת הסכימה. בוחרים אחת מהאפשרויות האלה:
מאגר סכימות: הסכימות שלכם מאוחסנות במאגר סכימות של Confluent. האפשרות הזו שימושית לשינוי סכימות ולניהול של כמה גרסאות. מוודאים שמאגר הסכימות נגיש לרשת של אשכול השירות המנוהל ל-Apache Kafka, ושהוא מתארח באותו אזור כמו העובדים של Dataflow. אפשר להשתמש במרשם סכימות גם בתרחישים של סכימה אחת וגם בתרחישים של כמה סכימות. מגדירים את ההגדרות הנוספות הבאות:
כתובת ה-URL של החיבור למאגר הסכימות: צריך לספק את כתובת ה-URL לחיבור למאגר הסכימות.
מצב אימות: אם המרשם דורש אימות, בוחרים באפשרות OAuth או TLS. אחרת, בוחרים באפשרות ללא.
קובץ סכימה יחיד: בוחרים באפשרות הזו אם כל ההודעות שלכם פועלות לפי סכימה קבועה אחת שמוגדרת בקובץ.
- קובץ של סכימת Avro לקובץ של אחסון בענן: הנתיב לקובץ של סכימת Avro שמשמש לפענוח כל ההודעות בנושא.
קידוד בינארי של Avro
אם בוחרים באפשרות הזו כפורמט ההודעה של Kafka, צריך להגדיר את ההגדרות הנוספות הבאות:
- קובץ של סכימת Avro לקובץ של אחסון בענן: הנתיב לקובץ של סכימת Avro שמשמש לפענוח כל ההודעות בנושא.
JSON
אם בוחרים באפשרות הזו בתור פורמט ההודעה של Kafka, לא נדרשות הגדרות נוספות.
ציון ההיסט של Kafka
כדי להימנע מעיבוד מחדש של הודעות כשצריך להפעיל מחדש עובדים בודדים או את כל צינור הנתונים, בוחרים באפשרות Commit offsets to Kafka. כך מובטח שצינור הנתונים ימשיך לעבד את הנתונים מהמקום שבו הוא הפסיק, ויימנע עיבוד כפול ואי-התאמות פוטנציאליות בנתונים.
בשדה Enter Consumer Group ID (הזנת מזהה קבוצת צרכנים), מזינים שם ייחודי לקבוצה של צינור הנתונים הזה. ברוב המקרים, רוצים שהצינור יקרא כל הודעה פעם אחת ושאפשר יהיה להפעיל אותו מחדש.
בשדה Default Kafka start offset, צינור הנתונים של Dataflow מציע שתי אפשרויות של היסט התחלתי. צריך לבחור אחת מהאפשרויות האלה:
המוקדם ביותר: מעבד הודעות מתחילת הנושא ב-Kafka.
האחרון: עיבוד ההודעות מתחיל מההיסט האחרון שזמין.
הגדרת יעד
האפשרויות האלה קובעות איך צינור עיבוד הנתונים כותב נתונים ל-Cloud Storage.
בשדה Destination, מזינים את הנתיב לקטגוריה וכוללים את הקידומת של שם הקובץ עבור קובצי הפלט. התחילית של הקובץ צריכה להסתיים בלוכסן. לדוגמה,
gs://test-bucket/test-prefix/בקטע משך החלון, מזינים את חלון הזמן לכתיבת נתונים ב-Cloud Storage. בוחרים את הפורמט המתאים (
Nsלשניות,Nmלדקות,Nhלשעות) בהתאם לדרישות עיבוד הנתונים.בשדה Output filename prefix of the files to write (קידומת לשם קובץ הפלט של הקבצים לכתיבה), אפשר לספק קידומת שתתווסף לכל קובץ פלט כדי לשפר את הארגון והזיהוי.
בשדה Maximum output shards (מספר מקסימלי של רסיסי פלט), מגדירים את המספר לאפס. אפשר לציין את מספר הרסיסים שייווצרו כשכותבים קבצים. הגדלת המספר יכולה להגדיל את התפוקה, אבל היא גם מובילה לעלויות גבוהות יותר בגלל עלויות גבוהות יותר של ערבוב נתונים. אם מגדירים את המספר לאפס, השירות בוחר מספר אופטימלי.
הגדרת תור של הודעות שלא ניתן להעביר
לפעמים אי אפשר לעבד הודעות בגלל נתונים פגומים, סוגי נתונים לא תואמים או אי התאמה בין סכימות.
כדי לטפל במקרים האלה, מפעילים את תור ההודעות שלא נמסרו בתבנית ומציינים שם של טבלה. התבנית יוצרת את הטבלה באמצעות סכימה סטנדרטית.
הגדרת הצפנה
כברירת מחדל, כל הנתונים במצב מנוחה ובזמן ההעברה מוצפנים על ידיGoogle-owned and Google-managed encryption key. אם יש לכם מפתחות הצפנה בניהול הלקוח (CMEK), אתם יכולים לבחור מפתחות משלכם. מידע נוסף על הגדרת CMEK זמין במאמר הגדרת הצפנת הודעות.
הגדרה של רשתות
צריך לציין את הרשת ואת רשת המשנה של האשכול בתבנית Dataflow. בקטע Optional parameters (פרמטרים אופציונליים) בתבנית אפשר להגדיר את הרשת של העובדים ב-Dataflow.
תבנית Kafka to Cloud Storage Dataflow מספקת כברירת מחדל עובדי Dataflow ברשת ברירת המחדל של הפרויקט. כדי לאפשר לאשכול של השירות המנוהל ל-Apache Kafka לשלוח נתונים ל-Cloud Storage דרך Dataflow, צריך לוודא שלעובדי Dataflow יש גישה לרשת של האשכול.
אם אשכול Kafka לא מחובר לרשת משנה ברשת ברירת המחדל של הפרויקט, מומלץ להשתמש ברשת ברירת המחדל של הפרויקט עבור אשכול Kafka.
מידע נוסף על הגדרת רשתות עם צינור Dataflow זמין במאמרים הבאים:
אם נתקלתם בבעיות בהגדרת הרשת של Dataflow, כדאי לעיין במדריך לפתרון בעיות ברשת של Dataflow.
הגדרת פרמטרים אופציונליים של Dataflow
מגדירים את הפרמטרים האופציונליים רק אם יודעים מה ההשפעה של ההגדרה על העובדים של Dataflow. הגדרות שגויות יכולות להשפיע על הביצועים או על העלות. הסברים מפורטים על כל אפשרות מופיעים במאמר פרמטרים אופציונליים.
מעקב
תבנית Dataflow לKafka ל-Cloud Storage מספקת חוויית מעקב שמאפשרת לכם לעיין ביומנים, במדדים ובשגיאות במסוף. חבילת כלי המעקב הזו זמינה כחלק מממשק המשתמש של Dataflow.
בכרטיסייה מדדים של משימות אפשר ליצור מרכזי בקרה בהתאמה אישית. בתבנית Dataflow Kafka to Cloud Storage, מומלץ להגדיר לוח בקרה של מדדי משימות כדי לעקוב אחרי המדדים הבאים:
קצב העברת נתונים: נפח הנתונים שעובר עיבוד בכל נקודת זמן. האפשרות הזו שימושית למעקב אחרי זרימת הנתונים דרך העבודה ולזיהוי בעיות פוטנציאליות בביצועים.
מידע נוסף זמין במאמר בנושא מעקב אחרי קצב העברת הנתונים ב-Dataflow.
עדכניות הנתונים: ההפרש בשניות בין חותמת הזמן של רכיב הנתונים לבין הזמן שבו האירוע עובר עיבוד בצינור הנתונים. כך אפשר לזהות צווארי בקבוק בביצועים ובמקורות הנתונים או ניסיונות חוזרים ונשנים.
מידע נוסף זמין במאמר בנושא מעקב אחרי עדכניות הנתונים ב-Dataflow.
Backlog: כמות הבייטים שממתינים לעיבוד. המידע הזה משמש לקבלת החלטות לגבי התאמה אוטומטית לעומס.
מידע נוסף על מעקב ב-Dataflow זמין במסמכי התיעוד בנושא מעקב ב-Dataflow.
פתרון בעיות
אם נתקלים בבעיות בביצועים של צינור Dataflow, Dataflow מספקת קבוצה מקיפה של כלים לפתרון בעיות ולאבחון.
ריכזנו כאן שני תרחישים נפוצים ומדריכים לפתרון בעיות שקשורים אליהם:
סקירה כללית על ניפוי באגים בצינורות עיבוד נתונים של Dataflow זמינה במאמר פתרון בעיות וניפוי באגים בצינורות עיבוד נתונים של Dataflow.
מגבלות ידועות
התבנית לא תומכת בהעברת פרטי כניסה לאימות ל-Schema Registry.
כשיוצרים את משימת Dataflow Kafka to Cloud Storage, מוודאים שהפרויקט Google Cloud מוגדר לאותו פרויקט שמכיל את האשכול של השירות המנוהל ל-Apache Kafka.