יצירת נתונים סינתטיים עבור אשכול של שירות מנוהל ל-Apache Kafka

במאמר הזה נסביר איך ליצור נתוני בדיקה סינתטיים עבור אשכול של שירות מנוהל של Google Cloud ל-Apache Kafka באמצעות Google Cloud המסוף.

במדריך הזה נעשה שימוש בתבנית Dataflow Streaming Data Generator כדי לפרסם באופן אוטומטי נתוני טלמטריה לדוגמה של משחק בנושא שירות מנוהל ל-Apache Kafka. הכלי ליצירת נתונים בסטרימינג הוא תבנית Dataflow שיוצרת רשומות בדיקה סינתטיות על סמך סכימה שצוינה, בקצב שניתן להגדרה. יצירת נתונים סינתטיים מאפשרת לכם לצפות בפעילות של אשכולות, לבדוק את הטיפול בעומס ולאמת מדדי מעקב בלי להתקין לקוח Kafka מקומי או לכתוב קוד מותאם אישית של יצרן. מידע נוסף על התבנית זמין במאמר תבנית Dataflow Streaming Data Generator.

לפני שמתחילים

לפני שמתחילים את המדריך הזה, צריך ליצור אשכול חדש של שירות מנוהל ל-Apache Kafka. אם כבר יש לכם אשכול, אתם יכולים לדלג על השלב הזה. מידע על התפקידים וההרשאות שנדרשים ליצירת אשכול זמין במאמר יצירה והצגה של אשכול. אם פועלים לפי המדריך הזה, צריך להשלים רק את הקטע יצירת אשכול לפני שחוזרים למדריך הזה.

איך יוצרים אשכול

המסוף

  1. עוברים אל הדף שירות מנוהל ל-Apache Kafka > Clusters.

    מעבר אל Clusters

  2. לוחצים על יצירה.
  3. בתיבה שם האשכול, מזינים שם לאשכול.
  4. ברשימה Region, בוחרים מיקום לאשכול.
  5. בקטע Network configuration (תצורת רשת), מגדירים את רשת המשנה שבה אפשר לגשת לאשכול:
    1. בקטע Project (פרויקט), בוחרים את הפרויקט.
    2. בקטע רשת, בוחרים את רשת ה-VPC.
    3. בקטע רשת משנה, בוחרים את רשת המשנה.
    4. לוחצים על סיום.
  6. לוחצים על יצירה.

אחרי שלוחצים על יצירה, מצב האשכול הוא Creating. כשהאשכול מוכן, המצב הוא Active.

gcloud

כדי ליצור אשכול Kafka, מריצים את הפקודה managed-kafka clusters create.

gcloud managed-kafka clusters create KAFKA_CLUSTER \
--location=REGION \
--cpu=3 \
--memory=3GiB \
--subnets=projects/PROJECT_ID/regions/REGION/subnetworks/SUBNET_NAME \
--async

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

  • KAFKA_CLUSTER: שם לאשכול Kafka
  • REGION: המיקום של האשכול
  • PROJECT_ID: מזהה הפרויקט
  • SUBNET_NAME: תת-הרשת שבה רוצים ליצור את האשכול, לדוגמה default

מידע על מיקומים נתמכים זמין במאמר בנושא מיקומים של שירות מנוהל ל-Apache Kafka.

הפקודה מופעלת באופן אסינכרוני ומחזירה מזהה פעולה:

Check operation [projects/PROJECT_ID/locations/REGION/operations/OPERATION_ID] for status.

כדי לעקוב אחרי ההתקדמות של פעולת היצירה, משתמשים בפקודה gcloud managed-kafka operations describe:

gcloud managed-kafka operations describe OPERATION_ID \
  --location=REGION

כשהאשכול מוכן, הפלט מהפקודה הזו כולל את הרשומה state: ACTIVE. מידע נוסף זמין במאמר בנושא מעקב אחרי פעולת יצירת האשכול.

התפקידים הנדרשים

כדי לקבל את ההרשאות שדרושות ליצירת נתונים סינתטיים עבור אשכול, צריך לבקש מהאדמין להקצות לכם את תפקידי ה-IAM הבאים בפרויקט:

להסבר על מתן תפקידים, ראו איך מנהלים את הגישה ברמת הפרויקט, התיקייה והארגון.

התפקידים המוגדרים מראש האלה מכילים את ההרשאות שנדרשות ליצירת נתונים סינתטיים לאשכול. כדי לראות בדיוק אילו הרשאות נדרשות, אפשר להרחיב את הקטע ההרשאות הנדרשות:

ההרשאות הנדרשות

כדי ליצור נתונים סינתטיים לאשכול, צריך את ההרשאות הבאות:

  • dataflow.jobs.create
  • dataflow.jobs.get
  • managedkafka.clusters.get
  • managedkafka.topics.get
  • managedkafka.topics.create
  • managedkafka.topics.publish
  • resourcemanager.projects.setIamPolicy

יכול להיות שתקבלו את ההרשאות האלה באמצעות תפקידים בהתאמה אישית או תפקידים מוגדרים מראש אחרים.

כדי לוודא שלחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine יש את ההרשאות שנדרשות להרצת משימת Dataflow, צריך לבקש מהאדמין להקצות לחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine את תפקידי ה-IAM הבאים בפרויקט:

להסבר על מתן תפקידים, ראו איך מנהלים את הגישה ברמת הפרויקט, התיקייה והארגון.

יכול להיות שהאדמין גם יוכל לתת לחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine את ההרשאות שנדרשות באמצעות תפקידים בהתאמה אישית או תפקידים מוגדרים מראש אחרים.

אם יש לכם הרשאות להקצאת תפקידי IAM,Google Cloud המסוף יציג בקשה להקצאת התפקידים הנדרשים במהלך תהליך יצירת הנתונים הסינתטיים. אם אין לכם הרשאה להקצות תפקידים, במסוף תוצג הודעה שבה תתבקשו לבקש מהאדמין להקצות את ההרשאות הנדרשות.

יצירת נתונים סינתטיים

כדי ליצור ולהפעיל את משימת Dataflow שמייצרת נתונים סינתטיים לנושא Kafka:

  1. במסוף Google Cloud , נכנסים לדף שירות מנוהל ל-Apache Kafka > Clusters.

    מעבר אל Clusters

  2. לוחצים על שם האשכול, למשל test-cluster.

  3. לוחצים על הכרטיסייה מקורות.

  4. בדף מקורות, בכרטיס יצירת נתונים סינתטיים, לוחצים על יצירת משימת Dataflow. החלונית Produce Data נפתחת.

  5. בחלונית Produce Data, בוחרים נושא מהרשימה הנפתחת Kafka topic, למשל test-topic. אם אין לכם נושא, אתם יכולים ליצור אחד:

    1. ברשימה הנפתחת Kafka topic, לוחצים על Create topic. החלונית Create topic תיפתח.
    2. בשדה Topic name, מזינים test-topic.
    3. משאירים את ערכי ברירת המחדל של Partition count (3) ושל Replication factor (3).
    4. לוחצים על יצירה.
  6. בשדה Output rate (QPS) (קצב הפלט (QPS)), מזינים את קצב השאילתות לשנייה שרוצים שהגנרטור ייצור, למשל 100. כך אפשר לבדוק איך האשכול מטפל בעומסים שונים.

  7. אם מופיעה אזהרה שלפיה לחשבון השירות של Dataflow אין את ההרשאות הנדרשות, לוחצים על Grant (הענקה) כדי להקצות את התפקידים הבאים:

    • Dataflow Worker‏ (roles/dataflow.worker)
    • Managed Kafka Client ‏ (roles/managedkafka.client)
  8. בחלונית Produce Data (יצירת נתונים), לוחצים על Create (יצירה) כדי להפעיל את משימת Dataflow.

    תוצג הודעה שהמשימה של Dataflow נוצרה.

  9. בהתראה, לוחצים על הצגת המשימה כדי לפתוח את הדף פרטי משימת Dataflow, שבו אפשר לראות את גרף המשימה, הסטטוס ומדדי הביצוע.

הצגת מדדים של אשכול

אחרי שהמשימה של Dataflow מתחילה, אפשר לראות את הנתונים הסינתטיים זורמים אל האשכול:

  1. בדף פרטי האשכול של test-cluster, לוחצים על הכרטיסייה מעקב.

  2. בודקים את התרשימים קצב העברת נתונים (בייט) ו5 הנושאים המובילים לפי נפח הנתונים שמועברים כדי לוודא שהנתונים מופקים באופן פעיל לנושא.

צפייה בהודעות

כדי לוודא שהודעות סינתטיות מתפרסמות בנושא, אפשר להשתמש באחת מהשיטות הבאות.

צפייה בכלי שורת הפקודה של Kafka

כדי לצרוך הודעות ישירות מהאשכול באמצעות כלי Kafka CLI במכונת VM של לקוח:

  1. מתחברים למכונה הווירטואלית של הלקוח באמצעות SSH. אם לא הגדרתם מכונה וירטואלית של לקוח, כדאי לעיין במאמר בנושא יצירת מכונה וירטואלית של לקוח.

  2. מקבלים את כתובת שרת ה-bootstrap של האשכול מ Google Cloud המסוף ומגדירים אותה כמשתנה סביבה ב-VM של הלקוח:

    1. במסוף Google Cloud , נכנסים לדף שירות מנוהל ל-Apache Kafka > Clusters.

      מעבר אל Clusters

    2. לוחצים על שם האשכול, למשל test-cluster.

    3. בדף פרטי האשכול, לוחצים על הגדרות.

    4. מעתיקים את הערך שמופיע בקטע Bootstrap URL.

  3. ב-VM של הלקוח, מגדירים את משתנה הסביבה:

    ```sh
    export BOOTSTRAP="BOOTSTRAP_URL"
    ```
    

    מחליפים את BOOTSTRAP_URL בכתובת ה-bootstrap שהעתקתם.

  4. מריצים את הפקודה kafka-console-consumer.sh כדי לקרוא את ההודעות:

    kafka-console-consumer.sh \
     --bootstrap-server $BOOTSTRAP \
     --topic TOPIC_ID \
     --from-beginning \
     --consumer.config client.properties
    

    מחליפים את TOPIC_ID בשם הנושא, למשל test-topic.

    רשומות נתוני המשחק הסינתטיים בסטרימינג מוצגות במסוף בזמן שהן נצרכות.

  5. מקישים על Ctrl+C כדי להפסיק את צריכת ההודעות.

הצגה ב-BigQuery

כדי להזרים נתונים מנושא Kafka ל-BigQuery ולצפות ברשומות:

  1. מכיוון שהנתונים הסינתטיים הם JSON גולמי, צריך ליצור באופן ידני את טבלת היעד ב-BigQuery לפני שיוצרים את המחבר. מידע על יצירת טבלה מופיע במאמר יצירת טבלה ריקה עם הגדרת סכימה. יוצרים טבלה בשם test-topic במערך הנתונים עם הסכימה הבאה:

    [
      {"name": "eventId", "type": "STRING"},
      {"name": "eventTimestamp", "type": "INTEGER"},
      {"name": "ipv4", "type": "STRING"},
      {"name": "ipv6", "type": "STRING"},
      {"name": "country", "type": "STRING"},
      {"name": "username", "type": "STRING"},
      {"name": "quest", "type": "STRING"},
      {"name": "score", "type": "INTEGER"},
      {"name": "completed", "type": "BOOLEAN"}
    ]
    
  2. יוצרים מחבר BigQuery Sink באשכול Connect כדי להזרים הודעות מהנושא לטבלה ב-BigQuery. כשמגדירים את המחבר, משתמשים במאפיינים לדוגמה הבאים ומחליפים את PROJECT_ID במזהה הפרויקט:

    bigQueryPartitionDecorator=false
    connector.class=com.wepay.kafka.connect.bigquery.BigQuerySinkConnector
    defaultDataset=test_dataset
    key.converter=org.apache.kafka.connect.storage.StringConverter
    project=PROJECT_ID
    tasks.max=3
    topics=test-topic
    value.converter=org.apache.kafka.connect.json.JsonConverter
    value.converter.schemas.enable=false
    
  3. אחרי שהמחבר מתחיל להזרים נתונים, עוברים לדף BigQuery במסוף Google Cloud .

    כניסה לדף BigQuery

  4. בחלונית Explorer, מרחיבים את מזהה הפרויקט ובוחרים את מערך הנתונים, test_dataset.

  5. לוחצים על שם הטבלה, test-topic.

  6. לוחצים על הכרטיסייה תצוגה מקדימה כדי לראות את הרשומות הסינתטיות שמוזרמות. לחלופין, לוחצים על Compose new query ומריצים את שאילתת ה-SQL הבאה:

    SELECT * FROM `PROJECT_ID.DATASET_ID.TABLE_ID` LIMIT 10;
    

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

    • PROJECT_ID: מזהה הפרויקט
    • DATASET_ID: מזהה מערך הנתונים, למשל test_dataset
    • TABLE_ID: מזהה הטבלה, למשל test-topic
  7. לוחצים על Run כדי לראות את רשומות הדוגמה בחלונית Query results.

    הערה: אל תשתמשו בשאילתת SELECT COUNT(*) כדי לאמת את הרשומות. המחבר משתמש ב-BigQuery Streaming API, ולכן הנתונים נכתבים בהתחלה למאגר זמני של נתונים זורמים. אפשר לראות את הנתונים מיד באמצעות SELECT *, אבל יכול לעבור כמה דקות עד שספירת השורות תתעדכן.

הסרת המשאבים

כדי לא לצבור חיובים לחשבון Google Cloud על המשאבים שבהם השתמשתם בדף הזה, פועלים לפי השלבים הבאים:

  1. נכנסים לדף Dataflow Jobs במסוף Google Cloud .

    מעבר אל 'משימות Dataflow'

  2. לוחצים על השם של המשימה שנוצרה לנושא.

  3. לוחצים על הפסקה.

  4. בוחרים באפשרות ביטול ואז לוחצים על הפסקת העבודה.

  5. אופציונלי: אם אתם לא צריכים יותר את אשכול Kafka, עוברים לדף Managed Service for Apache Kafka Clusters, בוחרים באפשרות test-cluster ולוחצים על Delete.

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