将 BigQuery 数据同步到 AlloyDB

本页介绍了如何将 BigQuery 中的表同步到 AlloyDB for PostgreSQL 实例中。

通过将 BigQuery 中的分析数据同步到 AlloyDB,您可以构建运营系统,该系统可受益于对数据湖的低延迟事务性访问。与就地查询数据的外部数据封装容器 (FDW) 不同,同步表会将数据移至 AlloyDB 存储空间,以实现最佳性能。

AlloyDB 提供了以下方法,可将 BigQuery 数据移至您的实例:

  • 一次性同步:创建 BigQuery 表的可写入独立副本。

  • 定期同步(镜像):创建只读本地表,该表会按预定时间自动刷新,例如每 6 小时或每天刷新一次。

性能和运营注意事项

使用 BigQuery 同步表时,请注意以下事项:

  • 资源用量:数据移动会消耗 CPU 和内存。对于非常大的表,请考虑在非高峰时段安排同步,以免影响主要事务性工作负载。
  • 数据可见性:在替换操作期间,系统会预先舍弃并重新创建现有目标表。在导入期间,查询最初会看到一个空表,随后随着批处理事务提交,新导入的数据会逐步显示。

准备工作

  1. 熟悉 bigquery_fdw 如何处理 BigQuery 数据类型和列映射,因为 alloydb_sync 扩展程序使用 bigquery_fdw 连接到 BigQuery。
  2. 登录您的 Google Cloud 账号。如果您是 Google Cloud新手,请 创建一个账号来评估我们的产品在实际场景中的表现。新客户还可获享 $300 赠金,用于运行、测试和部署工作负载。
  3. In the Google Cloud console, on the project selector page, select or create a Google Cloud project.

    Roles required to select or create a project

    • Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
    • Create a project: To create a project, you need the Project Creator role (roles/resourcemanager.projectCreator), which contains the resourcemanager.projects.create permission. Learn how to grant roles.

    Go to project selector

  4. Verify that billing is enabled for your Google Cloud project.

  5. In the Google Cloud console, on the project selector page, select or create a Google Cloud project.

    Roles required to select or create a project

    • Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
    • Create a project: To create a project, you need the Project Creator role (roles/resourcemanager.projectCreator), which contains the resourcemanager.projects.create permission. Learn how to grant roles.

    Go to project selector

  6. Verify that billing is enabled for your Google Cloud project.

  7. 启用创建和连接到 AlloyDB 所需的 Cloud API。

    启用 API

  8. 如需确认您要更改的项目的名称,请在确认项目步骤中点击下一步

  9. 启用 API 步骤中,点击启用以启用以下内容:

    • AlloyDB API
    • Compute Engine API
    • Cloud Resource Manager API
    • Service Networking API
    • BigQuery Storage API

    如果您计划使用与 AlloyDB 位于同一 Google Cloud 项目中的 VPC 网络配置与 AlloyDB 的网络连接,则需要使用 Service Networking API。

    如果您计划使用位于其他 Google Cloud 项目中的 VPC 网络配置与 AlloyDB 的网络连接,则需要使用 Compute Engine API 和 Cloud Resource Manager API。

  10. 确保您有一个现有的 BigQuery 表,可从中同步数据。如需了解详情,请参阅创建和使用 BigQuery 表

所需的角色

如需向 AlloyDB 集群服务账号授予对 BigQuery 数据集的访问权限,您需要以下权限:

  • BigQuery Data Viewer (roles/bigquery.dataViewer) 或具有 bigquery.tables.getbigquery.tables.getData 权限的任何自定义角色。针对服务账号授予此角色后,可提供从表或视图中读取数据和元数据的权限。
  • BigQuery Read Session User (roles/bigquery.readSessionUser) 或具有 bigquery.readsessions.createbigquery.readsessions.getData 权限的任何自定义角色。提供创建和使用读取会话的功能。
  • BigQuery Job User (roles/bigquery.jobUser) 或具有 bigquery.jobs.create 权限的任何自定义角色。提供创建和运行作业(包括查询作业)的功能。

配置扩展程序

在从 BigQuery 同步表之前,请启用所需的扩展程序并配置与 BigQuery 的连接。

  1. 创建 扩展程序:

    1. 按照将 psql 客户端连接到实例中的说明,使用 psql 客户端连接到 AlloyDB 实例。
    2. 运行以下命令:

      CREATE EXTENSION IF NOT EXISTS alloydb_sync;
      
  2. 如需让 AlloyDB 通过 BigQuery 进行身份验证,请创建用户映射。

    CREATE EXTENSION IF NOT EXISTS bigquery_fdw;
    CREATE SERVER IF NOT EXISTS BIGQUERY_SERVER_NAME FOREIGN DATA WRAPPER bigquery_fdw;
    CREATE USER MAPPING IF NOT EXISTS FOR USER SERVER BIGQUERY_SERVER_NAME;
    

    替换以下内容:

    • USER:数据库用户名或可访问 BigQuery 表的 IAM 用户。
    • BIGQUERY_SERVER_NAME:BigQuery 服务器的唯一标识符。在给定数据库中定义一次。 您可以将 BIGQUERY_SERVER_NAME 替换为您的服务器名称。

同步 BigQuery 表以进行一次性导出

您可以使用 psql 同步 BigQuery 表以进行一次性导出。

使用 psql 一次性同步 BigQuery 表

如需创建可修改的 BigQuery 数据副本,请使用 psql 运行 alloydb_sync.import_bq_table 函数。

SELECT alloydb_sync.import_bq_table(
  'PROJECT_ID.DATASET_ID.TABLE_ID',
  'ALLOYDB_DESTINATION_TABLE_NAME',
  'ON_EXISTS',
  ARRAY['PRIMARY_KEY_COLUMN']
);

替换以下内容:

  • PROJECT_ID:BigQuery 数据集所在项目的 ID。
  • DATASET_ID:相应表的 BigQuery 数据集的名称。对于具有 4 部分名称的 Iceberg 表,这是 Catalog.Namespace
  • TABLE_ID:BigQuery 表或视图的名称。
  • ALLOYDB_DESTINATION_TABLE_NAME:要在 AlloyDB 数据库中创建并导入数据的本地表的名称。您可以添加架构名称,例如 public.local_sales
  • ON_EXISTS:如果目标表已存在,则使用的策略。
  • PRIMARY_KEY_COLUMN:用作主键的可选列名称列表。

示例

以下示例展示了如何将名为 transactions 的表从 BigQuery 数据集同步到名为 public.local_sales 的新 AlloyDB 表:

SELECT alloydb_sync.import_bq_table(
    'my-gcp-project.sales_data.transactions',
    'public.local_sales',
    'replace'
);
on_exists 参数

on_exists 参数用于确定当目标表已存在于 AlloyDB 中时,该函数如何处理同步:

  • error:默认选项。如果目标表已存在,则停止同步。
  • skip:如果目标表已存在,则跳过同步。
  • replace:使用来自 BigQuery 的新数据替换现有的本地表。
主键支持

如果您以文本数组的形式提供可选的 primary_key 参数,AlloyDB 会创建以指定列为主键的表。

SELECT alloydb_sync.import_bq_table(
    'my-gcp-project.sales_data.transactions',
    'public.local_sales',
    ARRAY['transaction_id']
);

同步 BigQuery 表以进行定期导出

您可以使用 psql 同步 BigQuery 表以进行定期导出。

创建定期同步

如需维护与 BigQuery 数据保持同步的只读表,请使用 psql 运行 alloydb_sync.create_bq_sync_table 函数。

SELECT alloydb_sync.create_bq_sync_table(
    'PROJECT_ID.DATASET_ID.TABLE_ID',
    'ALLOYDB_DESTINATION_TABLE_NAME',
    'REFRESH_INTERVAL',
    'ON_EXISTS',
    ARRAY['PRIMARY_KEY_COLUMN']
);

替换以下内容:

  • PROJECT_ID.DATASET_ID.TABLE_ID:BigQuery 表或视图的完全限定名称,包括项目 ID、数据集 ID 和表 ID,以英文句点分隔。 对于具有 4 部分名称的 Iceberg 表,DATASET_ID 表示为 Catalog.Namespace。例如 my-gcp-project.sales_data.transactions
  • ALLOYDB_DESTINATION_TABLE_NAME:AlloyDB 数据库中要创建并同步数据到的本地表的名称。
  • REFRESH_INTERVAL:AlloyDB 定期刷新来自 BigQuery 的数据的间隔,例如 12 hours
  • ON_EXISTS:如果目标表已存在,则使用的策略。
  • PRIMARY_KEY_COLUMN:要用作主键的可选列名称列表。

示例

以下示例展示了如何创建每 12 小时刷新一次的客户资料镜像:

SELECT alloydb_sync.create_bq_sync_table(
    'my-gcp-project.crm_data.profiles',
    'public.customer_mirror',
    '12 hours',
    'replace'
);

监控和管理作业

发起同步后,您可以监控同步进度并管理作业。

检查作业状态

大型同步可能需要一些时间。您可以通过查询 job_status 视图来监控进度,包括已处理的记录数和预计完成时间:

SELECT
    import_id,
    status,
    records_processed,
    total_records,
    error
FROM alloydb_sync.job_status;

例如,如需取消作业,请运行以下命令:

SELECT alloydb_sync.cancel_import_job('85bb5dfa-dfb9-4017-9153-738f55abe4b1');

停止并删除同步作业

如需停止镜像 BigQuery 表并删除本地表,请使用 alloydb_sync.delete_bq_sync_table 函数:

SELECT alloydb_sync.delete_bq_sync_table('public.customer_mirror');

限制

从 BigQuery 同步表时,存在以下限制:

  • 此功能仅适用于 PostgreSQL 18 版。
  • 如果您DROP alloydb_sync 扩展程序,则必须先重启实例,然后才能再次创建该扩展程序。
  • 同步在事务中运行。如果导入作业中断或失败,系统会回滚导入的数据。
  • 如果两个用户同时启动同步作业,且目标表相同,则这两个表可能会相互覆盖。
  • 如果在新注册的同步表的初始后台导入期间发生任何中断,则该表在下一个预定的刷新间隔之前仍处于不完整状态。如需解决此问题,您可以使用 alloydb_sync.delete_bq_sync_table() 函数删除同步表,然后重新创建该表。
  • 不支持同步复杂的 BigQuery 类型,例如 ARRAYBYTESVECTORGEOGRAPHY。如需查看完整列表,请参阅受支持的 BigQuery 数据类型和列映射
  • 请勿手动删除复制的表。使用 alloydb_sync.delete_bq_sync_table() API 函数安全地舍弃表并刷新。
  • 如需删除使用 alloydb_sync 扩展程序的数据库,您必须使用 DROP DATABASE ... WITH (FORCE)
  • 如果 Postgres 数据库在导入运行时崩溃,元数据可能会卡在 RUNNING 状态,从而阻止未来的导入。您必须手动运行 UPDATE alloydb_sync.import_job_status SET status = 'FAILED' WHERE status = 'RUNNING'; 才能取消屏蔽。

价格

当您将数据从 BigQuery 同步到 AlloyDB 时,您需要按照 BigQuery 容量计算价格付费。

导出数据后,如果您将数据存储在 AlloyDB 中,则需要为此付费。如需了解详情,请参阅 AlloyDB for PostgreSQL 价格

后续步骤