使用跨云数据湖仓一体

Lakehouse for Apache Iceberg 支持通过跨云 Lakehouse 配置查询远程数据。配置完成后,系统支持在 BigQuery、Apache Spark 开源版本或 Managed Service for Apache Spark 中使用标准 SQL 访问数据。 除了分析查询之外,您还可以将联合数据用于 AI 驱动的数据洞见和治理:

  • 对话分析: 根据您的确切数据源(包括 跨云表)构建专用代理,以便通过一次 对话分析跨云数据。
  • Dataplex Catalog: 使用 Knowledge Catalog 功能,通过 联合数据源进行数据分析和分析洞见。

如需更深入地了解数据,您可以根据数据源(从项目、数据集和表到视图、图表和用户定义函数)编写专用代理。由于您的数据很少存储在一个位置,因此对话分析功能不仅适用于 BigQuery 标准表,还适用于 Lakehouse 管理的 Apache Iceberg 表以及 Databricks Unity、AWS Glue、SAP 和 Salesforce 等跨云 Lakehouse 来源。这样,您就可以打破数据孤岛,并通过一次对话分析跨云数据。

本页介绍了在设置跨云 Lakehouse 后如何查询远程数据。

准备工作

您必须先完成以下操作,然后才能查询数据:

  1. AWS GlueDatabricks Unity CatalogSnowflake 设置跨云 Lakehouse。
  2. 确保远程目录中有数据。

所需角色

如需获得查询联合数据所需的权限,请让管理员向您授予项目的以下 IAM 角色:

如需详细了解如何授予角色,请参阅管理对项目、文件夹和组织的访问权限

您也可以通过自定义角色或其他预定义角色来获取所需的权限。

查询数据

设置联合后,您可以在 BigQuery 中使用标准 SQL 或在 Managed Service for Apache Spark 中使用 Apache Spark 查询远程数据。

Lakehouse 会处理元数据转换和安全数据 访问,让您可以将远程 Apache Iceberg 表视为本地 环境中的表。 Google Cloud

从 BigQuery 查询

如需查询联合 Apache Iceberg 表,请使用标准 BigQuery SQL。表路径遵循 4 部分结构:project.federated_catalog.namespace.table。系统会自动处理缓存、凭据自动售卖和 CCI 传输路由。

SELECT
  user_id,
  action,
  COUNT(*) as total_actions
FROM `PROJECT_ID.FEDERATED_CATALOG_NAME.NAMESPACE_NAME.TABLE_NAME`
WHERE event_date >= '2026-04-01'
GROUP BY 1, 2;

替换以下内容:

  • PROJECT_ID:您的 Google Cloud 项目 ID。
  • FEDERATED_CATALOG_NAME:联合目录的名称。
  • NAMESPACE_NAME:目录中的命名空间。
  • TABLE_NAME:表的名称。
  • REGION:区域。 Google Cloud 例如 us-east4

您还可以使用 bq 命令行工具运行查询:

bq --location="REGION" --project_id="PROJECT_ID" query --use_legacy_sql=false \
  "SELECT * FROM \`PROJECT_ID.FEDERATED_CATALOG_NAME.NAMESPACE_NAME.TABLE_NAME\` LIMIT 10"

从 Managed Service for Apache Spark 查询

使用 X-Iceberg-Access-Delegation=vended-credentials启用 凭据自动售卖功能后,向 Managed Service for Apache Spark 提交PySpark 批量工作负载。Spark 将使用短期范围的自动售卖凭据安全地连接到 S3,而无需管理单独的 AWS 凭据或 S3 连接器。

  1. 为 Managed Service for Apache Spark 启用出站连接。

    Managed Service for Apache Spark 无法使用其默认 网络配置连接到 AWS S3。您必须预配 Cloud Router 和 Cloud NAT。

    gcloud compute routers create lakehouse-router \
      --network=NETWORK_NAME \
      --region=REGION
    
    gcloud compute routers nats create lakehouse-nat \
      --router=lakehouse-router \
      --auto-allocate-nat-external-ips \
      --nat-all-subnet-ip-ranges \
      --region=REGION

    替换以下内容:

    • NETWORK_NAME:Managed Service for Apache Spark 批量工作负载的网络(例如 default)。
    • REGION:Managed Service for Apache Spark 批量工作负载的区域。
  2. 创建 PySpark 应用文件并运行 PySpark 作业。

    from pyspark.sql import SparkSession
    spark = SparkSession.builder.appName("CATALOG_NAME").getOrCreate()
    
    df = spark.table("CATALOG_NAME.NAMESPACE_NAME.TABLE_NAME")
    df.show(10, truncate=False)

    将其上传到 Cloud Storage 中的 PYSPARK_FILE

    gcloud dataproc batches submit pyspark PYSPARK_FILE \
        --project=PROJECT_ID \
        --region=REGION \
        --version=RUNTIME_VERSION \
        --properties="\
        spark.sql.defaultCatalog=CATALOG_NAME,\
        spark.sql.catalog.CATALOG_NAME=org.apache.iceberg.spark.SparkCatalog,\
        spark.sql.catalog.CATALOG_NAME.type=rest,\
        spark.sql.catalog.CATALOG_NAME.uri=https://biglake.googleapis.com/iceberg/v1/restcatalog,\
        spark.sql.catalog.CATALOG_NAME.warehouse=bl://projects/PROJECT_ID/catalogs/FEDERATED_CATALOG_NAME,\
        spark.sql.catalog.CATALOG_NAME.header.x-goog-user-project=PROJECT_ID,\
        spark.sql.catalog.CATALOG_NAME.rest.auth.type=org.apache.iceberg.gcp.auth.GoogleAuthManager,\
        spark.sql.catalog.CATALOG_NAME.io-impl=IO_IMPL,\
        spark.sql.catalog.CATALOG_NAME.header.X-Iceberg-Access-Delegation=vended-credentials,\
        spark.sql.catalog.CATALOG_NAME.rest-metrics-reporting-enabled=false,\
        spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"

    替换以下内容:

    • NAMESPACE_NAME:联合目录中的命名空间。
    • TABLE_NAME:联合目录中的表名称。
    • CATALOG_NAME:本地 Spark 目录的名称(例如 my_catalog)。
    • PYSPARK_FILE:PySpark 应用文件的 gs:// Cloud Storage 路径。
    • REGION:Managed Service for Apache Spark 批量工作负载的区域。
    • RUNTIME_VERSION:Managed Service for Apache Spark 运行时版本,例如 2.3
    • PROJECT_ID:使用 Apache Iceberg REST 目录端点所产生的费用将计入该项目。
    • FEDERATED_CATALOG_NAME:联合目录的名称。
    • IO_IMPL:与底层存储匹配的 FileIO 实现。

    Spark 配置参数

    下表列出了所有连接所需的常见参数:

    参数 说明
    spark.sql.defaultCatalog 默认目录名称(例如 CATALOG_NAME)。
    spark.sql.catalog.CATALOG_NAME 目录实现类。设置为 org.apache.iceberg.spark.SparkCatalog
    spark.sql.catalog.CATALOG_NAME.type 目录后端类型。对于 Iceberg REST 目录,设置为 rest
    spark.sql.catalog.CATALOG_NAME.uri REST 目录端点的 URI。设置为 https://biglake.googleapis.com/iceberg/v1/restcatalog
    spark.sql.catalog.CATALOG_NAME.warehouse 联合目录的仓库位置路径。设置为 bl://projects/PROJECT_ID/catalogs/FEDERATED_CATALOG_NAME
    spark.sql.catalog.CATALOG_NAME.header.x-goog-user-project 用于结算和配额归因的 Google Cloud 项目 ID。设置为 PROJECT_ID
    spark.sql.extensions Iceberg SQL 语法和功能的 Spark 会话扩展。设置为 org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions

    下表列出了身份验证参数:

    参数 说明
    spark.sql.catalog.CATALOG_NAME.rest.auth.type 自定义身份验证管理器类。对于 OAuth 流程身份验证,设置为 org.apache.iceberg.gcp.auth.GoogleAuthManager
    spark.sql.catalog.CATALOG_NAME.oauth2-server-uri OAuth2 令牌服务器端点 URI。对于个人访问令牌 (PAT) 身份验证,设置为 https://oauth2.googleapis.com/token
    spark.sql.catalog.CATALOG_NAME.token Bearer 令牌或个人访问令牌 (PAT)。对于 PAT 身份验证,通常设置为 $(gcloud auth application-default print-access-token)

    下表列出了基于存储空间服务 (IO_IMPL) 的唯一参数:

    存储 spark.sql.catalog.CATALOG_NAME.io-impl 备注
    Amazon S3 org.apache.iceberg.aws.s3.S3FileIO 如果启用,则需要凭据自动售卖 (X-Iceberg-Access-Delegation=vended-credentials)。
    其他参数:spark.sql.catalog.CATALOG_NAME.s3.region(如需查看区域列表,请参阅 Amazon S3 端点和配额)。
    Google Cloud Storage org.apache.iceberg.gcp.gcs.GCSFileIO 其他参数:spark.sql.catalog.CATALOG_NAME.gcs.oauth2.refresh-credentials-endpoint=https://oauth2.googleapis.com/token
    Azure Blob Storage org.apache.iceberg.azure.adlsv2.ADLSFileIO 无需其他存储参数。

    对于 Snowflake,查询 STRING 列时可能会遇到问题,因为这些列会自动进行存储优化。您可以通过以下两种方式之一解决此问题:

    • 选项 1:在 Spark 中停用矢量化:将以下 Spark 配置属性添加到 --properties 标志:
    • spark.sql.iceberg.vectorization.enabled=false
    • spark.sql.catalog.CATALOG_NAME.table-override.read.parquet.vectorization.enabled=false

    • 选项 2:更改 Snowflake 中的序列化政策:将 Snowflake 中表的存储序列化政策更改为 COMPATIBLE

后续步骤