Managed Airflow (דור 3) | Managed Airflow (דור 2) | Managed Airflow (דור 1 מדור קודם)
בדף הזה מוסבר איך ליצור ולשנות טבלאות BigLake Iceberg ב-BigQuery באמצעות אופרטורים של Airflow בסביבת Managed Airflow.
מידע על טבלאות BigLake Iceberg ב-BigQuery
טבלאות BigLake Iceberg ב-BigQuery מספקות את הבסיס לבניית אגמי נתונים בפורמט פתוח ב- Google Cloud. טבלאות BigLake Iceberg ב-BigQuery מציעות את אותו ניהול מלא כמו טבלאות BigQuery רגילות, אבל הנתונים מאוחסנים בדלי אחסון בבעלות הלקוח. טבלאות BigLake Iceberg ב-BigQuery תומכות בפורמט הפתוח של טבלאות Iceberg, כדי לשפר את יכולת הפעולה ההדדית עם מנועי חישוב בקוד פתוח ושל צד שלישי, על עותק יחיד של נתונים.
לפני שמתחילים
חשוב לוודא שאתם מכירים את המגבלות על טבלאות BigLake Iceberg ב-BigQuery. אותן מגבלות חלות כשעובדים עם טבלאות BigLake Iceberg באמצעות אופרטורים של Airflow.
יוצרים קישור למשאבים ב-Cloud לקטגוריה של Cloud Storage שבה תמוקם טבלת Iceberg של BigLake.
חשוב לוודא שההרשאות ב-IAM מוקצות באופן הבא:
לחשבון השירות של הסביבה צריכים להיות תפקידי IAM לעבודה עם טבלאות BigLake Iceberg ב-BigQuery. פועלים לפי ההוראות שמפורטות במאמר טבלאות BigLake ל-Apache Iceberg ב-BigQuery.
לחשבון השירות שמשויך לקישור למשאבים ב-Cloud צריכים להיות תפקידי IAM לקריאה ולכתיבה של נתונים ב-Cloud Storage. פועלים לפי ההוראות שמפורטות במאמר טבלאות BigLake ל-Apache Iceberg ב-BigQuery.
יצירת טבלת BigLake Iceberg ב-BigQuery
כדי ליצור טבלת Iceberg ב-BigLake ב-BigQuery, משתמשים ב-BigQueryCreateTableOperator באותו אופן כמו בטבלאות אחרות ב-BigQuery. בשדה biglakeConfiguration, מציינים את ההגדרה של הטבלה.
import datetime
from airflow.models.dag import DAG
from airflow.providers.google.cloud.operators.bigquery import BigQueryCreateTableOperator
with DAG(
"bq_iceberg_dag",
start_date=datetime.datetime(2025, 1, 1),
schedule=None,
) as dag:
create_iceberg_table = BigQueryCreateTableOperator(
task_id="create_iceberg_table",
project_id="PROJECT_ID",
dataset_id="DATASET_ID",
table_id="TABLE_NAME",
table_resource={
"schema": {
"fields": [
{"name": "order_id", "type": "INTEGER", "mode": "REQUIRED"},
{"name": "customer_id", "type": "INTEGER", "mode": "REQUIRED"},
{"name": "amount", "type": "INTEGER", "mode": "REQUIRED"},
{"name": "created_at", "type": "TIMESTAMP", "mode": "REQUIRED"},
]
},
"biglakeConfiguration": {
"connectionId": "CONNECTION_NAME",
"storageUri": "STORAGE_URI",
"fileFormat": "PARQUET",
"tableFormat": "ICEBERG",
}
}
)
מחליפים את מה שכתוב בשדות הבאים:
-
PROJECT_ID: מזהה הפרויקט. DATASET_ID: מערך נתונים קיים.-
TABLE_NAME: השם של הטבלה שיוצרים. -
CONNECTION_NAME: השם של הקישור למשאבים ב-Cloud בפורמטprojects/PROJECT_ID/locations/REGION/connections/CONNECTION_ID. -
STORAGE_URI: URI של Cloud Storage מוגדר במלואו עבור הטבלה. לדוגמה:gs://example-bucket/iceberg-table.
שליחת שאילתות לטבלאות BigLake Iceberg ב-BigQuery
אחרי שיוצרים טבלת Iceberg ב-BigLake, אפשר להריץ עליה שאילתות באמצעות BigQueryInsertJobOperator כרגיל. האופרטור לא צריך הגדרות נוספות ספציפיות לטבלאות BigLake Iceberg.
import datetime
from airflow.models.dag import DAG
from airflow.providers.google.cloud.operators.bigquery import BigQueryInsertJobOperator
with DAG(
"bq_iceberg_dag_query",
start_date=datetime.datetime(2025, 1, 1),
schedule=None,
) as dag:
insert_values = BigQueryInsertJobOperator(
task_id="iceberg_insert_values",
configuration={
"query": {
"query": f"""
INSERT INTO `TABLE_ID` (order_id, customer_id, amount, created_at)
VALUES
(101, 19, 1, TIMESTAMP '2025-09-15 10:15:00+00'),
(102, 35, 2, TIMESTAMP '2025-09-14 10:15:00+00'),
(103, 36, 3, TIMESTAMP '2025-09-12 10:15:00+00'),
(104, 37, 4, TIMESTAMP '2025-09-11 10:15:00+00')
""",
"useLegacySql": False,
}
}
)
מחליפים את מה שכתוב בשדות הבאים:
-
TABLE_IDבמזהה הטבלה, בפורמטPROJECT_ID.DATASET_ID.TABLE_NAME.