יצירת Lakehouse באמצעות Spark וקטלוג זמן הריצה של Lakehouse
ארכיטקטורת lakehouse משלבת את הגמישות של אגם נתונים עם תכונות ניהול הנתונים של מחסן נתונים. במאמר הזה מוסבר איך להגדיר lakehouse ב- Google Cloud. אתם משתמשים ב-Apache Iceberg כפורמט הטבלה, ב-Managed Service for Apache Spark לעיבוד ובקטלוג זמן הריצה של Lakehouse, Iceberg REST Catalog, לניהול מאוחד של מטא-נתונים.
הארכיטקטורה הזו משתמשת בפורמטים פתוחים של טבלאות כמו Iceberg כדי להוסיף לנתונים ב-Cloud Storage יכולות של מחסן נתונים, כמו טרנזקציות והתפתחות סכימה. הגישה הזו יוצרת מקור מידע אמין אחד לנתונים שלכם, שאפשר לגשת אליו ממנועים שונים.
לפני שמתחילים
- נכנסים לחשבון Google Cloud . אם אתם משתמשים חדשים ב- Google Cloud, צרו חשבון כדי שתוכלו להעריך את הביצועים של המוצרים שלנו בתרחישים מהעולם האמיתי. לקוחות חדשים מקבלים בחינם גם קרדיט בשווי 300$ להרצה, לבדיקה ולפריסה של עומסי העבודה.
-
In the Google Cloud console, on the project selector page, select or create a Google Cloud project.
Roles required to select or create a project
- Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
-
Create a project: To create a project, you need the Project Creator role
(
roles/resourcemanager.projectCreator), which contains theresourcemanager.projects.createpermission. Learn how to grant roles.
-
Verify that you have the permissions required to complete this guide.
-
Verify that billing is enabled for your Google Cloud project.
Enable the Dataproc, BigQuery, and Cloud Storage APIs.
Roles required to enable APIs
To enable APIs, you need the
serviceusage.services.enablepermission. If you created the project, then you likely already have this permission through the Owner role (roles/owner). Otherwise, you can get this permission through the Service Usage Admin role (roles/serviceusage.serviceUsageAdmin). Learn how to grant roles.-
In the Google Cloud console, on the project selector page, select or create a Google Cloud project.
Roles required to select or create a project
- Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
-
Create a project: To create a project, you need the Project Creator role
(
roles/resourcemanager.projectCreator), which contains theresourcemanager.projects.createpermission. Learn how to grant roles.
-
Verify that you have the permissions required to complete this guide.
-
Verify that billing is enabled for your Google Cloud project.
Enable the Dataproc, BigQuery, and Cloud Storage APIs.
Roles required to enable APIs
To enable APIs, you need the
serviceusage.services.enablepermission. If you created the project, then you likely already have this permission through the Owner role (roles/owner). Otherwise, you can get this permission through the Service Usage Admin role (roles/serviceusage.serviceUsageAdmin). Learn how to grant roles.- יצירת קטגוריה של Cloud Storage לאחסון נתוני Iceberg.
התפקידים הנדרשים
כדי להריץ את הדוגמאות בדף הזה, צריך תפקידים מסוימים בממשק לניהול הזהויות והרשאות הגישה (IAM). יכול להיות שהתפקידים האלה כבר הוקצו, בהתאם למדיניות הארגון. כדי לבדוק את התפקידים שהוקצו, ראו האם צריך להקצות תפקידים?.
כדי לקרוא הסבר על מתן תפקידים, ראו איך מנהלים את הגישה ברמת הפרויקט, התיקייה והארגון.
תפקידי משתמשים
כדי לקבל את ההרשאות שדרושות ליצירת אשכול של Managed Service for Apache Spark, צריך לבקש מהאדמין להקצות לכם את תפקידי ה-IAM הבאים:
- עריכת פרויקטים ב-Dataproc (
roles/dataproc.editor) בפרויקט - משתמש בחשבון שירות (
roles/iam.serviceAccountUser) בחשבון השירות של Compute Engine שמוגדר כברירת מחדל
תפקיד בחשבון שירות
כדי לוודא שלחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine יש את ההרשאות שנדרשות ליצירת אשכול של Managed Service for Apache Spark, צריך לבקש מהאדמין להקצות לחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine את תפקיד ה-IAM Dataproc Worker (roles/dataproc.worker) בפרויקט.
יצירת אשכול Managed Service for Apache Spark
יוצרים אשכול של Managed Service for Apache Spark עם הרכיבים האופציונליים Iceberg ו-Jupyter.
כדי ליצור את האשכול, מריצים את הפקודה
gcloudהבאה:gcloud dataproc clusters create CLUSTER_NAME \ --project=PROJECT_ID \ --region=REGION \ --image-version=2.3-debian12 \ --optional-components=ICEBERG,JUPYTER \ --enable-component-gateway \ --properties 'dataproc:dataproc.lineage.enabled=true'מחליפים את מה שכתוב בשדות הבאים:
-
CLUSTER_NAME: שם לאשכול. -
PROJECT_ID: מזהה הפרויקט ב- Google Cloud . -
REGION: האזור Google Cloud של האשכול, לדוגמה,us-central1.
הערה: לא צריך להגדיר את
dataproc:dataproc.lineage.enabled=trueכדי שהקטלוג של זמן הריצה של Lakehouse, קטלוג REST של Iceberg, יפעל בצורה תקינה. הוא מתווסף למעקב אחר מקורות נתונים בדוגמה למעקב אחר מקורות נתונים שבהמשך.-
מתחברים לאשכול באמצעות מחברת Jupyter. אפשר להשתמש ב-notebook של Vertex AI Workbench או להפעיל notebook ישירות באשכול.
הגדרת סשן Spark
ב-Jupyter Notebook, יוצרים סשן Spark שמוגדר לשימוש בקטלוג של Lakehouse runtime Iceberg REST Catalog.
import pyspark
from pyspark.context import SparkContext
from pyspark.sql import SparkSession
catalog_name = "CATALOG_NAME"
spark = SparkSession.builder.appName("APP_NAME") \
.config(f'spark.sql.catalog.{catalog_name}', 'org.apache.iceberg.spark.SparkCatalog') \
.config(f'spark.sql.catalog.{catalog_name}.type', 'rest') \
.config(f'spark.sql.catalog.{catalog_name}.uri', 'https://biglake.googleapis.com/iceberg/v1beta/restcatalog') \
.config(f'spark.sql.catalog.{catalog_name}.warehouse', 'gs://GCS_BUCKET') \
.config(f'spark.sql.catalog.{catalog_name}.header.x-goog-user-project', 'PROJECT_ID') \
.config(f'spark.sql.catalog.{catalog_name}.rest.auth.type', 'org.apache.iceberg.gcp.auth.GoogleAuthManager') \
.config(f'spark.sql.catalog.{catalog_name}.io-impl', 'org.apache.iceberg.gcp.gcs.GCSFileIO') \
.config(f'spark.sql.catalog.{catalog_name}.rest-metrics-reporting-enabled', 'false') \
.config('spark.sql.extensions', 'org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions') \
.config('spark.sql.defaultCatalog', catalog_name) \
.getOrCreate()
מחליפים את מה שכתוב בשדות הבאים:
-
CATALOG_NAME: שם לקטלוג Iceberg, למשלiceberg_catalog. -
APP_NAME: השם של אפליקציית Spark. -
GCS_BUCKET: הקטגוריה ב-Cloud Storage שבה יישמרו נתוני טבלת Iceberg. -
PROJECT_ID: מזהה הפרויקט ב- Google Cloud .
ניהול נתונים באמצעות Spark SQL
אחרי שמגדירים את סשן Spark, משתמשים ב-Spark SQL כדי לבצע פעולות של ניהול נתונים.
יוצרים מרחב שמות. בקטלוג של Lakehouse בזמן הריצה, קטלוג REST של Iceberg, מרחב שמות תואם למערך נתונים ב-BigQuery.
spark.sql("CREATE NAMESPACE IF NOT EXISTS NAMESPACE_NAME") spark.sql("USE NAMESPACE_NAME")מחליפים את
NAMESPACE_NAMEבשם של מרחב השמות, לדוגמה,spark_lakehouse.יוצרים טבלת בסיס בפורמט Iceberg ומכניסים נתונים.
spark.sql("DROP TABLE IF EXISTS base_table PURGE") spark.sql("CREATE TABLE base_table (id LONG) USING iceberg") spark.sql("INSERT INTO base_table VALUES 0, 1, 2, 3, 4") spark.sql("SELECT * FROM base_table").show()הפלט אמור להיראות כך:
+---+ | id| +---+ | 0| | 1| | 2| | 3| | 4| +---+יוצרים טבלה שנייה לנתונים חדשים.
spark.sql("DROP TABLE IF EXISTS newdata PURGE") spark.sql("CREATE TABLE newdata(id LONG) USING iceberg") spark.sql("INSERT INTO newdata VALUES 3, 4, 5, 6") spark.sql("SELECT * FROM newdata").show()הפלט אמור להיראות כך:
+---+ | id| +---+ | 3| | 4| | 5| | 6| +---+ממזגים את הנתונים החדשים לטבלת הבסיס.
spark.sql("""MERGE INTO base_table USING newdata ON base_table.id = newdata.id WHEN MATCHED THEN UPDATE SET base_table.id = newdata.id WHEN NOT MATCHED THEN INSERT * """) spark.sql("SELECT * FROM base_table").show()הפלט אמור להיראות כך:
+---+ | id| +---+ | 0| | 1| | 2| | 3| | 4| | 5| | 6| +---+עדכון רשומות בטבלת הבסיס.
spark.sql( "UPDATE base_table SET id = (id + 100) WHERE (id % 2 == 0)" ) spark.sql("SELECT * FROM base_table").show()הפלט אמור להיראות כך:
+---+ | id| +---+ | 3| |104| | 5| |106| |100| |102| | 1| +---+מחיקת רשומות מטבלת הבסיס.
spark.sql("DELETE FROM base_table WHERE (id % 2 == 0)") spark.sql("SELECT * FROM base_table").show()הפלט אמור להיראות כך:
+---+ | id| +---+ | 3| | 5| | 1| +---+
שאילתה על תמונת מצב היסטורית
אפשר לאחזר גרסה קודמת של טבלה על ידי שליחת שאילתה למזהה ספציפי של תמונת מצב. הפעולה הזו נקראת גם 'מסע בזמן'.
מאחזרים את מזהה תמונת המצב של גרסת הטבלה לפני הפעולות
MERGE,UPDATEו-DELETE.snapshot_ids = spark.sql( "SELECT snapshot_id FROM `NAMESPACE_NAME`.`base_table`.snapshots" ).collect() oldest_snapshot_id = snapshot_ids[1]["snapshot_id"]מחליפים את
NAMESPACE_NAMEבמרחב השמות שיצרתם.מריצים שאילתה על הטבלה באמצעות מזהה התמונה ששוחזר.
df = ( spark.read.format("iceberg") .option("versionAsOf", oldest_snapshot_id) .load("base_table") ) df.show()בפלט מוצג המצב של הטבלה אחרי פעולת
MERGEאבל לפני פעולותUPDATEאוDELETE.+---+ | id| +---+ | 0| | 1| | 2| | 3| | 4| | 5| | 6| +---+
גילוי שושלת נתונים
אפשר לעקוב אחרי תנועת הנתונים בין טבלאות של קטלוג Iceberg REST של Lakehouse בזמן ריצה באמצעות שיוך נתונים למקור, שזמין בגרסה 2.2 ואילך של תמונות של Managed Service for Apache Spark.
דוגמה לשושלת נתונים
יוצרים טבלאות Iceberg של מקור ויעד, ואז מעתיקים את הנתונים.
spark.sql("DROP TABLE IF EXISTS source_table PURGE") spark.sql("DROP TABLE IF EXISTS target_table PURGE") spark.sql("CREATE TABLE source_table (id LONG) USING iceberg") spark.sql("""CREATE TABLE target_table USING ICEBERG AS SELECT max(id) as top_id FROM source_table """)נכנסים לדף Search ב-Knowledge Catalog במסוף Google Cloud .
מחפשים אחת מהטבלאות ולוחצים על הכרטיסייה
Lineage:
דוגמה לגרף של השתלשלות הנתונים בדף Knowledge Catalog במסוף Google Cloud . במעקב אחר מקורות נתונים מזוהים גם הייצוגים הלוגיים (טבלת קטלוג זמן הריצה של Lakehouse) וגם הייצוגים הפיזיים (Cloud Storage) של טבלאות קטלוג Iceberg REST של זמן הריצה של Lakehouse.
בעיה מוכרת בשושלת הנתונים
יכול להיות שבחלק מהאשכולות של Managed Service for Apache Spark לא ייווצר שושלת נתונים מלאה בגלל בעיה בספריית OpenLineage.
פתרון עקיף: בהגדרות של סשן Spark, מגדירים את המאפיין spark.sql.catalog.{catalog_name}.uri לערך https://biglake.googleapis.com/iceberg/v1beta/restcatalog.
המאמרים הבאים
- מידע נוסף על קטלוג REST של Iceberg בזמן הריצה של Lakehouse
- לומדים על התכונות של Apache Iceberg.
- איך שולחים שאילתות לנתוני Iceberg מקטלוג של Lakehouse runtime
- מידע נוסף על שושלת נתונים ו-Managed Service for Apache Spark