개방형 경계 없는 레이크하우스에서 교차 클라우드 속성 그래프 빌드

다음 튜토리얼에서는 데이터를 이동하지 않고 개방형 경계 없는 레이크하우스와 Apache Iceberg REST 카탈로그 엔드포인트를 사용하여 두 개의 서로 다른 클라우드에서 데이터 사일로를 통합하는 단일 BigQuery 그래프를 빌드하는 방법을 보여줍니다.

시작하기 전에

시작하기 전에 환경을 설정하고 필요한 API를 사용 설정하세요.

  1. 프로젝트 및 리전을 설정하고 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
    
  2. 로더의 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개를 로드합니다.

  1. 버킷 및 카탈로그를 만듭니다.

    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"
    
  2. 다음 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"))
    
  3. 스크립트를 실행하여 테이블을 로드합니다.

    python load_gcs.py
    
  4. 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에 표시합니다.

  1. 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);
    
  2. 서비스 주체에 읽기 액세스 권한을 부여합니다.

    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>`;
    
  3. 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=-
    
  4. 제휴 카탈로그를 만듭니다.

    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
    
  5. 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"
    
  6. 동기화를 확인합니다.

    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'
    

교차 클라우드 속성 그래프 만들기

노드 및 에지 테이블이 두 클라우드에 걸쳐 있는 그래프를 만듭니다.

  1. 리전에 데이터 세트를 만듭니다.

    bq --location="$REGION" mk --dataset "${PROJECT_ID}:xcloud_graph"
    
  2. 그래프를 만듭니다.

    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
    )
    

교차 클라우드 그래프 쿼리

그래프를 탐색하여 두 클라우드 간의 관계를 쿼리합니다.

  1. 두 클라우드에 있는 관계를 탐색하는 쿼리를 실행합니다.

    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;
    
  2. 각 고객에게 제품을 공급하는 국가를 찾는 쿼리를 실행합니다.

    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 계정에 청구되지 않도록 하려면 다음 단계를 수행합니다.

요금이 계속 청구되지 않도록 생성한 리소스를 삭제합니다.

  1. 그래프 데이터 세트를 삭제합니다.

    bq rm -r -f -d "${PROJECT_ID}:xcloud_graph"
    
  2. 제휴 카탈로그를 삭제합니다.

    gcloud alpha biglake iceberg catalogs delete dbx_fed_catalog --project="$PROJECT_ID"
    
  3. 다음 스크립트를 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)
    
  4. 스크립트를 실행합니다.

    python drop_gcs.py
    
  5. Cloud Storage 카탈로그를 삭제합니다.

    gcloud biglake iceberg catalogs delete "$GCS_CATALOG" --project="$PROJECT_ID"
    
  6. 보안 비밀을 삭제합니다.

    export CLOUDSDK_API_ENDPOINT_OVERRIDES_SECRETMANAGER="https://secretmanager.${REGION}.rep.googleapis.com/"
    gcloud secrets delete dbx-fed-sp --project="$PROJECT_ID" --location="$REGION"
    
  7. 버킷을 삭제합니다.

    gcloud storage rm -r "$GCS_BUCKET"
    

다음 단계