יצירת גרף נכסים חוצה-עננים על Lakehouse פתוח ללא גבולות
במדריך הבא מוסבר איך לבנות גרף BigQuery יחיד שמאחד נתונים מבודדים בשני עננים שונים באמצעות Lakehouse פתוח ללא גבולות ונקודות קצה של קטלוג Apache Iceberg REST בלי להעביר נתונים.
לפני שמתחילים
לפני שמתחילים, צריך להגדיר את הסביבה ולהפעיל את ממשקי ה-API הנדרשים.
מגדירים את הפרויקט והאזור ומפעילים את ממשקי ה-API:
export PROJECT_ID="your-gcp-project-id" export REGION="us-east4" gcloud config set project "$PROJECT_ID" gcloud services enable \ biglake.googleapis.com \ bigquery.googleapis.com \ secretmanager.googleapis.com \ storage.googleapis.comיוצרים סביבה וירטואלית של Python לטוען:
python3 -m venv iceberg-venv source iceberg-venv/bin/activate pip install --quiet "pyiceberg[pyarrow]"
יצירת Google Cloud הסתעפות
מגדירים קטלוג REST פתוח של Apache Iceberg שמגובה על ידי קטגוריה של Cloud Storage, וטוענים לתוכו שלושה טבלאות Iceberg.
יוצרים את הקטגוריה ואת הקטלוג:
export GCS_BUCKET="gs://${PROJECT_ID}-xcloud-lake" export GCS_CATALOG="gcs_lake" gcloud storage buckets create "$GCS_BUCKET" \ --project="$PROJECT_ID" \ --location="$REGION" gcloud biglake iceberg catalogs create "$GCS_CATALOG" \ --project="$PROJECT_ID" \ --catalog-type=biglake \ --primary-location="$REGION" \ --default-location="$GCS_BUCKET"שומרים את סקריפט Python הבא בשם
load_gcs.pyכדי לאכלס את הטבלאות:import subprocess, pyarrow as pa from pyiceberg.catalog.rest import RestCatalog from pyiceberg.schema import Schema from pyiceberg.types import NestedField, StringType, LongType, DoubleType import os PROJECT = os.environ["PROJECT_ID"] CATALOG = os.environ["GCS_CATALOG"] TOKEN = subprocess.check_output( ["gcloud", "auth", "application-default", "print-access-token"], text=True ).strip() cat = RestCatalog( name=CATALOG, uri="https://biglake.googleapis.com/iceberg/v1/restcatalog", warehouse=f"bl://projects/{PROJECT}/catalogs/{CATALOG}", token=TOKEN, **{"header.x-goog-user-project": PROJECT, "header.X-Iceberg-Access-Delegation": "vended-credentials"}, ) cat.create_namespace_if_not_exists("retail") def mk(name, schema, table): ident = ("retail", name) try: cat.drop_table(ident) except Exception: pass t = cat.create_table(ident, schema=schema) t.append(table) print(f" {name}: {table.num_rows} rows") # customers mk("customers", Schema(NestedField(1, "customer_id", StringType()), NestedField(2, "name", StringType()), NestedField(3, "region", StringType())), pa.table({ "customer_id": ["C1", "C2", "C3", "C4", "C5", "C6"], "name": ["Ana", "Ben", "Cara", "Dan", "Eve", "Finn"], "region": ["west", "west", "east", "east", "west", "south"], })) # orders mk("orders", Schema(NestedField(1, "order_id", StringType()), NestedField(2, "customer_id", StringType()), NestedField(3, "status", StringType()), NestedField(4, "amount", DoubleType())), pa.table({ "order_id": ["O1","O2","O3","O4","O5","O6","O7","O8","O9","O10"], "customer_id": ["C1","C1","C2","C3","C3","C4","C5","C5","C6","C2"], "status": ["shipped"]*8 + ["pending","shipped"], "amount": [156.0,89.0,120.0,147.0,89.0,199.0,25.0,88.0,80.0,224.0], })) # order_items mk("order_items", Schema(NestedField(1, "order_item_id", StringType()), NestedField(2, "order_id", StringType()), NestedField(3, "product_id", StringType()), NestedField(4, "quantity", LongType()), NestedField(5, "amount", DoubleType())), pa.table({ "order_item_id": [f"OI{i}" for i in range(1, 16)], "order_id": ["O1","O1","O2","O3","O3","O4","O5","O6","O6","O7","O8","O9","O10","O10","O2"], "product_id": ["P1","P2","P3","P1","P5","P4","P3","P6","P8","P7","P1","P2","P4","P5","P6"], "quantity": [1,2,1,1,3,1,1,1,4,2,1,1,1,2,1], "amount": [120.0,36.0,89.0,120.0,27.0,199.0,89.0,25.0,88.0,80.0,120.0,18.0,199.0,18.0,25.0], })) print("tables:", cat.list_tables("retail"))מריצים את הסקריפט כדי לטעון את הטבלאות:
python load_gcs.pyבודקים את הטבלאות מ-BigQuery:
bq --location="$REGION" query --use_legacy_sql=false \ 'SELECT customer_id, name, region FROM `'"$PROJECT_ID"'.gcs_lake.retail.customers` ORDER BY customer_id'
יצירת רכזת AWS
הצגת Databricks Unity Catalog ב-BigQuery דרך Lakehouse פתוח ללא גבולות.
בעורך Databricks SQL, יוצרים את הטבלאות products ו-suppliers:
CREATE SCHEMA IF NOT EXISTS `<CATALOG>`.retail; CREATE OR REPLACE TABLE `<CATALOG>`.retail.suppliers ( supplier_id STRING, name STRING, country STRING ) USING ICEBERG; INSERT INTO `<CATALOG>`.retail.suppliers VALUES ('S1','Acme','USA'), ('S2','Globex','Germany'), ('S3','Initech','Japan'), ('S4','Umbrella','UK'); CREATE OR REPLACE TABLE `<CATALOG>`.retail.products ( product_id STRING, name STRING, category STRING, supplier_id STRING, price DOUBLE ) USING ICEBERG; INSERT INTO `<CATALOG>`.retail.products VALUES ('P1','Widget','Gadgets', 'S1',120.0), ('P2','Gizmo','Gadgets', 'S1', 18.0), ('P3','Sprocket','Parts', 'S2', 89.0), ('P4','Cog','Parts', 'S2',199.0), ('P5','Bolt','Parts', 'S3', 9.0), ('P6','Nut','Parts', 'S3', 25.0), ('P7','Gear','Machinery', 'S4', 40.0), ('P8','Axle','Machinery', 'S4', 22.0);הקצאת גישת קריאה לחשבון השירות:
GRANT USE CATALOG ON CATALOG `<CATALOG>` TO `<SP_APPLICATION_ID>`; GRANT USE SCHEMA, SELECT, EXTERNAL USE SCHEMA ON SCHEMA `<CATALOG>`.retail TO `<SP_APPLICATION_ID>`;אחסון פרטי הכניסה ב-Secret Manager:
export CLOUDSDK_API_ENDPOINT_OVERRIDES_SECRETMANAGER="https://secretmanager.${REGION}.rep.googleapis.com/" printf '{"client_id":"<SP_CLIENT_ID>","client_secret":"<SP_CLIENT_SECRET>"}' \ | gcloud secrets create dbx-fed-sp \ --project="$PROJECT_ID" \ --location="$REGION" \ --data-file=-יוצרים את הקטלוג המאוחד:
export DBX_HOST="<your-workspace-host-without-https>" export UC_CATALOG="<CATALOG>" gcloud alpha biglake iceberg catalogs create dbx_fed_catalog \ --project="$PROJECT_ID" \ --catalog-type=federated \ --federated-catalog-type=unity \ --unity-instance-name="$DBX_HOST" \ --unity-catalog-name="$UC_CATALOG" \ --secret-name="projects/${PROJECT_ID}/locations/${REGION}/secrets/dbx-fed-sp" \ --primary-location="$REGION" \ --refresh-interval=330sמעניקים לחשבון השירות של BigLake גישה לסוד:
export BLIRC_SA="<paste-the-blirc-...-service-account-from-the-output>" export CLOUDSDK_API_ENDPOINT_OVERRIDES_SECRETMANAGER="https://secretmanager.${REGION}.rep.googleapis.com/" gcloud secrets add-iam-policy-binding dbx-fed-sp \ --project="$PROJECT_ID" \ --location="$REGION" \ --member="serviceAccount:${BLIRC_SA}" \ --role="roles/secretmanager.secretAccessor"מאמתים את הסנכרון:
bq --location="$REGION" ls dbx_fed_catalog.retail bq --location="$REGION" query --use_legacy_sql=false \ 'SELECT product_id, name, category, supplier_id, price FROM `'"$PROJECT_ID"'.dbx_fed_catalog.retail.products` ORDER BY product_id'
יצירת גרף הנכסים בין העננים
יוצרים גרף שבו טבלאות הצמתים והקשתות משתרעות על שני העננים.
כדי ליצור מערך נתונים באזור שלכם:
bq --location="$REGION" mk --dataset "${PROJECT_ID}:xcloud_graph"יוצרים את הגרף:
CREATE OR REPLACE PROPERTY GRAPH `PROJECT_ID.xcloud_graph.retail_xcloud` NODE TABLES ( `PROJECT_ID.gcs_lake.retail.customers` AS Customer KEY (customer_id) LABEL Customer PROPERTIES (customer_id, name, region), `PROJECT_ID.gcs_lake.retail.orders` AS Orders KEY (order_id) LABEL OrderNode PROPERTIES (order_id, status, amount), `PROJECT_ID.gcs_lake.retail.order_items` AS OrderItems KEY (order_item_id) LABEL OrderItem PROPERTIES (order_item_id, quantity, amount), `PROJECT_ID.dbx_fed_catalog.retail.products` AS Products KEY (product_id) LABEL Product PROPERTIES (product_id, name, category, price), `PROJECT_ID.dbx_fed_catalog.retail.suppliers` AS Suppliers KEY (supplier_id) LABEL Supplier PROPERTIES (supplier_id, name, country) ) EDGE TABLES ( `PROJECT_ID.gcs_lake.retail.orders` AS Placed KEY (order_id) SOURCE KEY (customer_id) REFERENCES Customer (customer_id) DESTINATION KEY (order_id) REFERENCES Orders (order_id) LABEL PLACED, `PROJECT_ID.gcs_lake.retail.order_items` AS Contains_Item KEY (order_item_id) SOURCE KEY (order_id) REFERENCES Orders (order_id) DESTINATION KEY (order_item_id) REFERENCES OrderItems (order_item_id) LABEL CONTAINS_ITEM, `PROJECT_ID.gcs_lake.retail.order_items` AS Is_Product KEY (order_item_id) SOURCE KEY (order_item_id) REFERENCES OrderItems (order_item_id) DESTINATION KEY (product_id) REFERENCES Products (product_id) LABEL IS_PRODUCT, `PROJECT_ID.dbx_fed_catalog.retail.products` AS Supplied_By KEY (product_id) SOURCE KEY (product_id) REFERENCES Products (product_id) DESTINATION KEY (supplier_id) REFERENCES Suppliers (supplier_id) LABEL SUPPLIED_BY )
שאילתה של הגרף בין העננים
אפשר לעבור בין הצמתים בגרף כדי לשאול שאילתות לגבי הקשרים בין העננים.
מריצים שאילתה שסורקת קשרים שנמצאים בשני העננים:
GRAPH `PROJECT_ID.xcloud_graph.retail_xcloud` MATCH (c:Customer)-[:PLACED]-> (o:OrderNode)-[:CONTAINS_ITEM]-> (li:OrderItem)-[:IS_PRODUCT]-> (p:Product)-[:SUPPLIED_BY]-> (s:Supplier) RETURN c.name AS customer, o.order_id AS order_id, p.name AS product, p.category AS category, s.name AS supplier, s.country AS supplier_country, li.quantity AS qty ORDER BY customer, order_id LIMIT 50;מריצים שאילתה כדי לגלות באילו מדינות נמצאים הספקים של כל לקוח:
GRAPH `PROJECT_ID.xcloud_graph.retail_xcloud` MATCH (c:Customer)-[:PLACED]-> (:OrderNode)-[:CONTAINS_ITEM]-> (:OrderItem)-[:IS_PRODUCT]-> (:Product)-[:SUPPLIED_BY]-> (s:Supplier) RETURN DISTINCT c.name AS customer, s.country AS supplier_country ORDER BY customer, supplier_country;
הסרת המשאבים
כדי לא לצבור חיובים לחשבון Google Cloud על המשאבים שבהם השתמשתם בדף הזה, פועלים לפי השלבים הבאים:
כדי להימנע מחיובים שוטפים, צריך להסיר את המשאבים שיצרתם.
מוחקים את מערך הנתונים של הגרף:
bq rm -r -f -d "${PROJECT_ID}:xcloud_graph"מחיקת הקטלוג המאוחד:
gcloud alpha biglake iceberg catalogs delete dbx_fed_catalog --project="$PROJECT_ID"שומרים את הסקריפט הבא בשם
drop_gcs.pyכדי להסיר את טבלאות GCS ואת מרחב השמות:import subprocess, os from pyiceberg.catalog.rest import RestCatalog PROJECT = os.environ["PROJECT_ID"] CATALOG = os.environ["GCS_CATALOG"] TOKEN = subprocess.check_output( ["gcloud", "auth", "application-default", "print-access-token"], text=True ).strip() cat = RestCatalog( name=CATALOG, uri="https://biglake.googleapis.com/iceberg/v1/restcatalog", warehouse=f"bl://projects/{PROJECT}/catalogs/{CATALOG}", token=TOKEN, **{"header.x-goog-user-project": PROJECT}, ) for tbl in cat.list_tables("retail"): cat.drop_table(tbl) cat.drop_namespace("retail") print("emptied", CATALOG)מריצים את הסקריפט:
python drop_gcs.pyמחיקת הקטלוג של Cloud Storage:
gcloud biglake iceberg catalogs delete "$GCS_CATALOG" --project="$PROJECT_ID"מחיקת הסוד:
export CLOUDSDK_API_ENDPOINT_OVERRIDES_SECRETMANAGER="https://secretmanager.${REGION}.rep.googleapis.com/" gcloud secrets delete dbx-fed-sp --project="$PROJECT_ID" --location="$REGION"מוחקים את הקטגוריה:
gcloud storage rm -r "$GCS_BUCKET"