개방형 경계 없는 레이크하우스에서 교차 클라우드 속성 그래프 빌드
다음 튜토리얼에서는 데이터를 이동하지 않고 개방형 경계 없는 레이크하우스와 Apache Iceberg REST 카탈로그 엔드포인트를 사용하여 두 개의 서로 다른 클라우드에서 데이터 사일로를 통합하는 단일 BigQuery 그래프를 빌드하는 방법을 보여줍니다.
시작하기 전에
시작하기 전에 환경을 설정하고 필요한 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 스포크 만들기
Cloud Storage 버킷으로 지원되는 개방형 Apache Iceberg REST 카탈로그를 설정하고 여기에 Iceberg 테이블 3개를 로드합니다.
버킷 및 카탈로그를 만듭니다.
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.pyBigQuery에서 테이블을 확인합니다.
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에 표시합니다.
Databricks SQL 편집기에서 제품 및 공급업체 테이블을 만듭니다.
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=330sBigLake 서비스 계정에 보안 비밀에 대한 액세스 권한을 부여합니다.
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.pyCloud 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"
다음 단계
- BigQuery 그래프에 대해 자세히 알아봅니다.
- 경계 없는 레이크하우스 개념에 대해 자세히 알아봅니다.
- 경계 없는 레이크하우스를 사용하는 방법을 알아봅니다.
- GQL 쿼리 문을 참조합니다.