ייבוא קובצי Parquet מאחסון ל-Lakehouse ללא גבולות באמצעות Dataflow

אתם יכולים להשתמש בתוכנית בסיסית ליצירת משימות Dataflow כדי להוסיף קובצי Apache Parquet קיימים מאחסון מבוסס-ענן (Cloud Storage,‏ Amazon S3 או Azure Blob Storage) לטבלת Apache Iceberg ב-Lakehouse ללא גבולות.

התהליך הזה משתמש בטרנספורמציה IcebergAddFiles. אם קובצי Parquet נמצאים ב-Cloud Storage, הטרנספורמציה הזו רושמת את הקבצים ב-Lakehouse בלי להעביר או לכתוב מחדש את הנתונים הבסיסיים. אם הקבצים שלכם נמצאים במערכת אחסון חיצונית כמו Amazon S3, הם מועתקים ל-Cloud Storage כדי לאפשר שאילתות מהירות יותר דרך Lakehouse, ואז הם נרשמים.

אפשר להשתמש בפרטי החיבור הבאים כדי להוסיף קובצי Parquet מאחסון מבוסס-ענן לטבלת Apache Iceberg ב-Lakehouse.

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

  1. מפעילים את ממשקי ה-API של Dataflow,‏ BigQuery ו-Lakehouse.

  2. כדי לקבל את ההרשאות שדרושות ליצירת המשאבים, צריך לבקש מהאדמין להקצות לכם את התפקידים הנדרשים בניהול הזהויות והרשאות הגישה (IAM) בפרויקט.

  3. יוצרים קטלוג, מרחב שמות וטבלה של Lakehouse כדי לייבא אליהם נתונים.

  4. יוצרים קטגוריה לאחסון בענן (Cloud Storage,‏ Amazon S3 או Azure Blob Storage) ומעלים אליה את קובצי ה-Parquet.

  5. אם קטגוריית האחסון מבוססת-הענן שבה אתם משתמשים היא לא Cloud Storage של Google, אתם צריכים ליצור קטגוריה של Cloud Storage כדי לאחסן את יומני השגיאות של העבודה.

תמיכה ומגבלות

יש מגבלות על ייבוא קובצי Parquet באחסון מבוסס-ענן ל-Lakehouse באמצעות Dataflow:

  • נתוני המקור צריכים להיות בפורמט Apache Parquet ומאוחסנים ב-Cloud Storage, ב-Amazon S3 או ב-Azure Blob Storage.
  • התכונה הזו תומכת רק בצינורות להעברת נתונים באצווה.

ייבוא קובצי Parquet ל-Lakehouse

כדי לייבא קובצי Parquet מאחסון מבוסס-ענן לטבלת Iceberg ב-Lakehouse, אפשר להשתמש בממשק המשתמש של הכלי ליצירת משימות Dataflow.

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

    מעבר אל Lakehouse

  2. בוחרים את הקטלוג, מרחב השמות והטבלה שאליהם רוצים לייבא נתונים.

  3. בדף פרטי הטבלה, לוחצים על ייבוא טבלה ואז בוחרים באפשרות מקובצי Apache Parquet (ייבוא מספר פריטים).

    ייפתח הדף Job builder של Dataflow.

  4. בקטע מקורות:

    1. פותחים את רשומת המקור CreateGlobalInput שכבר נוצרה.

    2. בקטע העריכה YAML source configuration, מזינים נתיב אחד או יותר לקובצי Parquet ברצף elements.

      כדי לשפר את יעילות הייבוא, כשרושמים מספר גדול של קבצים, כדאי לציין כמה קבוצות של קבצים (globs). לדוגמה:

      reshuffle: true
      elements:
        -   gs://BUCKET_NAME/restaurant-data/2023/*.parquet
        -   gs://BUCKET_NAME/restaurant-data/2024/*.parquet
      
    3. לוחצים על סיום.

  5. בקטע Transforms (טרנספורמציות):

    1. לוחצים על קטע הטרנספורמציה IcebergAddFiles כדי לפתוח אותו.

    2. בשדה Iceberg table (טבלת Iceberg), מזינים את מרחב השמות ואת שם הטבלה. לדוגמה: NAMESPACE .TABLE_NAME .

    3. בקטע מאפייני הקטלוג, מגדירים את הפריטים הבאים:

      1. warehouse: המיקום של הקטלוג ב-Cloud Storage. לדוגמה, gs://CATALOG_PATH.

      2. header.x-goog-user-project: מזהה הפרויקט ב- Google Cloud : PROJECT_ID.

      3. לוחצים על סיום.

    4. אם אתם מעבירים נתונים מ-Amazon S3 או מ-Azure Blob Storage, אתם צריכים לספק הגדרות נוספות כדי להעתיק קובצי Parquet ל-Cloud Storage. הפעולה הזו לא נדרשת אם הקבצים כבר נמצאים ב-Cloud Storage.

      1. לוחצים על הקטע של טרנספורמציית CopyFilesToGCS כדי לפתוח אותו.

      2. מגדירים את הערך של פרמטר ההגדרה gcs_file_path כדי לציין את הקטגוריה המוגדרת במלואה של Cloud Storage שאליה רוצים להעתיק את הקבצים הזמניים. מומלץ להשתמש באותה קטגוריה של Cloud Storage שבה משתמש מחסן הנתונים של Lakehouse.

      3. לוחצים על סיום.

      ‫1 לוחצים על הקטע Dataflow Options (אפשרויות של Dataflow) כדי לפתוח אותו.

      1. אם קובצי Parquet נמצאים ב-Amazon S3, לוחצים על add additional pipeline options (הוספת אפשרויות נוספות של צינור) כדי לספק ל-Apache Beam pipeline options related to S3 (אפשרויות של צינור שקשורות ל-S3). לדוגמה, s3_region_name,‏ s3_access_key_id,‏ s3_secret_access_key והערכים התואמים שלהם.

      2. אם קובצי ה-Parquet שלכם נמצאים ב-Azure Blob Storage, לוחצים על add additional pipeline options כדי לספק pipeline options related to Azure של Apache Beam. לדוגמה, azure_connection_string,‏ blob_service_endpoint,‏ azure_managed_identity_client_id והערכים התואמים שלהם.

  6. בקטע יעדים:

    1. לוחצים על יעד Write results כדי לפתוח אותו.

    2. בשדה JSON location (מיקום קובץ ה-JSON), מציינים את המיקום ב-Cloud Storage ואת שם הקובץ שבו ייכתבו תוצאות השגיאות. לדוגמה:

      gs://BUCKET_NAME/errors/errors.json
      
    3. לוחצים על סיום.

  7. בקטע Dataflow Options (אפשרויות של Dataflow), לוחצים על Run job (הפעלת משימה).

אם אתם צריכים להתאים אישית עוד יותר את צינור עיבוד הנתונים של Dataflow שמשמש לרישום קובצי Parquet, אתם יכולים לעשות זאת באמצעות טופס ליצירת משימות או עורך YAML.

בדיקת הפלט של המשימה

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

  1. ברשימת משימות Dataflow, בודקים שהסטטוס של המשימה הוא Succeeded.

    מעבר לדף Jobs

  2. אם המשימה נכשלת או שיש בה שגיאות, אפשר לבדוק את קובץ היומן של שגיאות ה-JSON ב-Cloud Storage כדי לקבל פרטים.

    כניסה לדף Buckets

  3. במסוף Google Cloud , עוברים לדף Studio ב-BigQuery.

    כניסה ל-BigQuery

  4. מזינים שאילתת SQL בעורך השאילתות כדי לבדוק את הטבלה. אפשר להשתמש במוסכמות של PROJECT_ID.CATALOG>NAMESPACE.TABLE_NAME כדי לשלוח שאילתות.

    SELECT * FROM `PROJECT_ID.CATALOG>NAMESPACE.TABLE_NAME` LIMIT 10
    
  5. לוחצים על הפעלה.

  6. בודקים את תוצאות השאילתה כדי לוודא שהנתונים עובדו בצורה נכונה.

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