使用无边界湖仓一体

Lakehouse for Apache Iceberg 支持通过无边界湖仓一体配置查询远程数据。配置完成后,系统支持在 BigQuery 中使用标准 SQL、Apache Spark 的开源版本或 Managed Service for Apache Spark 进行数据访问。除了分析查询之外,您还可以将联合数据用于 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 表一样处理远程 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 提交 PySpark 批量工作负载,并启用 credential vending。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 Catalog 端点所产生的费用将计入该项目。
    • 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 自定义身份验证管理器类。设置为 org.apache.iceberg.gcp.auth.GoogleAuthManager 以进行 OAuth 流程身份验证。
    spark.sql.catalog.CATALOG_NAME.oauth2-server-uri OAuth2 令牌服务器端点 URI。设置为 https://oauth2.googleapis.com/token 以进行个人访问令牌 (PAT) 身份验证。
    spark.sql.catalog.CATALOG_NAME.token 不记名令牌或个人访问令牌 (PAT)。通常设置为 $(gcloud auth application-default print-access-token) 以进行 PAT 身份验证。

    下表列出了根据您的存储空间服务 (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

后续步骤