在開放式無邊界 Lakehouse 上建構跨雲端屬性圖

下列教學課程說明如何建構單一 BigQuery 圖表,使用開放式無邊界 Lakehouse 和 Apache Iceberg REST 目錄端點,統一兩個不同雲端中的資料孤島,而不需移動資料。

事前準備

開始之前,請先設定環境並啟用必要的 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 bucket 支援的開放式 Apache Iceberg REST 目錄,並將三個 Iceberg 資料表載入其中。

  1. 建立 bucket 和目錄:

    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