使用 Dataflow 将存储空间中的 Parquet 文件导入 Lakehouse 运行时目录

您可以使用 Dataflow 作业构建器蓝图将基于云的存储空间(Cloud Storage、Amazon S3 或 Azure Blob Storage)中的现有 Apache Parquet 文件添加到无边界 Lakehouse 中的 Apache Iceberg 表。

此过程使用 IcebergAddFiles 转换。如果您的 Parquet 文件位于 Cloud Storage 中,此转换会将这些文件注册到 Lakehouse,而无需移动或重写底层数据。如果您的文件位于 Amazon S3 等外部存储系统中,则会先复制到 Cloud Storage 以便通过 Lakehouse 更快地查询,然后再进行注册。

使用以下连接详细信息将基于云的存储空间中的 Parquet 文件添加到 Lakehouse 中的 Apache Iceberg 表。

准备工作

  1. 启用 Dataflow、BigQuery 和 Lakehouse API。

  2. 如需获得创建资源所需的权限,请让管理员向您授予项目的必要 Identity and Access Management (IAM) 角色。

  3. 创建 Lakehouse 目录、命名空间和表,以便将数据导入其中。

  4. 创建一个基于云的存储桶(Cloud Storage、Amazon S3 或 Azure Blob 存储),并将 Parquet 文件上传到该存储桶。

  5. 如果您使用的云端存储桶不是 Google 的 Cloud Storage,请创建一个 Cloud Storage 存储桶来存储作业错误日志。

支持和限制

使用 Dataflow 将云端存储中的 Parquet 文件导入 Lakehouse 时,存在以下限制:

  • 源数据必须采用 Apache Parquet 格式,并存储在 Cloud Storage、Amazon S3 或 Azure Blob Storage 中。
  • 此功能仅支持批处理流水线。

将 Parquet 文件导入 Lakehouse

按照以下步骤操作,使用 Dataflow 作业构建器界面将 Parquet 文件从基于云的存储空间导入到 Lakehouse 中的 Iceberg 表。

  1. 在 Google Cloud 控制台中,前往 Lakehouse 页面。

    前往 Lakehouse

  2. 选择要将数据导入到的目录、命名空间和表。

  3. 表详细信息页面上,点击 导入表,然后选择从 Apache Parquet 文件导入(批量)

    随即会打开 Dataflow 作业构建器页面。

  4. 来源部分中:

    1. 打开已创建的 CreateGlobalInput 源条目。

    2. YAML 源配置编辑器部分中,以 elements 序列输入 Parquet 文件的一个或多个路径。

      为了提高导入效率,在注册大量文件时,请指定多组文件 (glob)。例如:

      reshuffle: true
      elements:
        -   gs://BUCKET_NAME/restaurant-data/2023/*.parquet
        -   gs://BUCKET_NAME/restaurant-data/2024/*.parquet
      
    3. 点击完成

  5. 转换部分中:

    1. 点击 IcebergAddFiles 转换部分以将其打开。

    2. Iceberg 表字段中,输入命名空间和表名称。例如:NAMESPACE TABLE_NAME

    3. 目录属性下,配置以下项:

      1. warehouse:目录的 Cloud Storage 位置。 例如 gs://CATALOG_PATH

      2. header.x-goog-user-project:您的 Google Cloud 项目 ID:PROJECT_ID

      3. 点击完成

    4. 如果您要从 Amazon S3 或 Azure Blob Storage 迁移,则需要提供额外的配置才能将 Parquet 文件复制到 Cloud Storage。如果您的文件已在 Cloud Storage 中,则无需执行此操作。

      1. 点击 CopyFilesToGCS 转换部分以将其打开。

      2. 设置 gcs_file_path 配置参数的值,以提供要将临时文件复制到的完全限定的 Cloud Storage 存储桶。建议使用 Lakehouse 仓储所用的同一 Cloud Storage 存储桶。

      3. 点击完成

      4. 点击 Dataflow 选项部分以将其打开。

      5. 如果您的 Parquet 文件位于 Amazon S3 中,请点击添加其他流水线选项,以提供与 S3 相关的 Apache Beam 流水线选项。 例如,s3_region_names3_access_key_ids3_secret_access_key 及其对应的值。

      6. 如果您的 Parquet 文件位于 Azure Blob Storage 中,请点击添加其他流水线选项,以提供与 Azure 相关的 Apache Beam 流水线选项。例如,azure_connection_stringblob_service_endpointazure_managed_identity_client_id 及其对应的值。

  6. 接收器部分中:

    1. 点击 Write results sink 以将其打开。

    2. JSON 位置字段中,指定用于写入错误结果的 Cloud Storage 位置和文件名。例如:

      gs://BUCKET_NAME/errors/errors.json
      
    3. 点击完成

  7. Dataflow 选项部分中,点击运行作业

如果您需要进一步自定义用于注册 Parquet 文件的 Dataflow 流水线,可以使用作业构建器表单或 YAML 编辑器来完成此操作。

检查作业输出

作业完成后,您可以在 BigQuery 中查询 Iceberg 表,验证数据是否已注册到该表中。

  1. 在 Dataflow 作业列表中,检查作业状态是否为成功

    转到作业

  2. 如果作业失败或出现错误,请检查 Cloud Storage 中的 JSON 错误日志文件以了解详情。

    进入“存储桶”

  3. 在 Google Cloud 控制台中,前往 BigQuery Studio 页面。

    转到 BigQuery

  4. 在查询编辑器中,输入 SQL 查询以检查表。您可以使用 PROJECT_ID.CATALOG.NAMESPACE.TABLE_NAME 惯例进行查询。

    SELECT * FROM `PROJECT_ID.CATALOG.NAMESPACE.TABLE_NAME` LIMIT 10
    
  5. 点击 运行

  6. 查看查询结果,确保数据已正确处理。

后续步骤