本页介绍了如何将 BigQuery 中的表同步到 AlloyDB for PostgreSQL 实例中。
通过将 BigQuery 中的分析数据同步到 AlloyDB,您可以构建运营系统,该系统可受益于对数据湖的低延迟事务性访问。与就地查询数据的外部数据封装容器 (FDW) 不同,同步表会将数据移至 AlloyDB 存储空间,以实现最佳性能。
AlloyDB 提供了以下方法,可将 BigQuery 数据移至您的实例:
一次性同步:创建 BigQuery 表的可写入独立副本。
定期同步(镜像):创建只读本地表,该表会按预定时间自动刷新,例如每 6 小时或每天刷新一次。
性能和运营注意事项
使用 BigQuery 同步表时,请注意以下事项:
- 资源用量:数据移动会消耗 CPU 和内存。对于非常大的表,请考虑在非高峰时段安排同步,以免影响主要事务性工作负载。
- 数据可见性:在替换操作期间,系统会预先舍弃并重新创建现有目标表。在导入期间,查询最初会看到一个空表,随后随着批处理事务提交,新导入的数据会逐步显示。
准备工作
- 熟悉
bigquery_fdw如何处理 BigQuery 数据类型和列映射,因为alloydb_sync扩展程序使用bigquery_fdw连接到 BigQuery。 -
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 theresourcemanager.projects.createpermission. Learn how to grant roles.
-
Verify that billing is enabled for your Google Cloud project.
-
启用创建和连接到 AlloyDB 所需的 Cloud API。
如需确认您要更改的项目的名称,请在确认项目步骤中点击下一步。
在启用 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。
- 确保您有一个现有的 BigQuery 表,可从中同步数据。如需了解详情,请参阅创建和使用 BigQuery 表。
所需的角色
如需向 AlloyDB 集群服务账号授予对 BigQuery 数据集的访问权限,您需要以下权限:
- BigQuery Data Viewer (
roles/bigquery.dataViewer) 或具有bigquery.tables.get和bigquery.tables.getData权限的任何自定义角色。针对服务账号授予此角色后,可提供从表或视图中读取数据和元数据的权限。 - BigQuery Read Session User (
roles/bigquery.readSessionUser) 或具有bigquery.readsessions.create和bigquery.readsessions.getData权限的任何自定义角色。提供创建和使用读取会话的功能。 - BigQuery Job User (
roles/bigquery.jobUser) 或具有bigquery.jobs.create权限的任何自定义角色。提供创建和运行作业(包括查询作业)的功能。
配置扩展程序
在从 BigQuery 同步表之前,请启用所需的扩展程序并配置与 BigQuery 的连接。
创建 扩展程序:
- 按照将 psql 客户端连接到实例中的说明,使用 psql 客户端连接到 AlloyDB 实例。
运行以下命令:
CREATE EXTENSION IF NOT EXISTS alloydb_sync;
如需让 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 版。
- 如果您
DROPalloydb_sync扩展程序,则必须先重启实例,然后才能再次创建该扩展程序。 - 同步在事务中运行。如果导入作业中断或失败,系统会回滚导入的数据。
- 如果两个用户同时启动同步作业,且目标表相同,则这两个表可能会相互覆盖。
- 如果在新注册的同步表的初始后台导入期间发生任何中断,则该表在下一个预定的刷新间隔之前仍处于不完整状态。如需解决此问题,您可以使用
alloydb_sync.delete_bq_sync_table()函数删除同步表,然后重新创建该表。 - 不支持同步复杂的 BigQuery 类型,例如
ARRAY、BYTES、VECTOR和GEOGRAPHY。如需查看完整列表,请参阅受支持的 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 价格。
后续步骤
- 详细了解将数据导出到 AlloyDB(反向 ETL)。
- 详细了解如何使用 FDW 就地查询 BigQuery 数据。