从 PostgreSQL 实时同步
本文介绍如何通过 Flink CDC pipeline 将 PostgreSQL 的数据实时同步至 StarRocks。Pipeline 先复制每张表中已有的数据,然后在秒级内持续应用 PostgreSQL 中的每一条 INSERT、UPDATE 和 DELETE。Pipeline 会自动创建 StarRocks 表,因此无需单独同步库表结构。
导入操作需要目标表的 INSERT 权限。如果您的用户账号没有 INSERT 权限,请参考 GRANT 给用户赋权,语法为 GRANT INSERT ON TABLE <table_name> IN DATABASE <database_name> TO { ROLE <role_name> | USER <user_identity>}。
基本原理
Flink CDC pipeline 是一个 YAML 文件,包含 source(PostgreSQL)、sink(StarRocks),以及可选的 route 和 transform 规则。将该文件提交到 Flink 集群后,它会作为一个 Flink 作业运行:
- PostgreSQL source 为每张匹配的表创建快照,然后通过逻辑复制槽(replication slot)从 PostgreSQL 预写日志(WAL)中读取变更。
- StarRocks sink 为 StarRocks 中尚不存在的每张源表创建一张主键表,然后通过 Stream Load 导入数据。由于每张表都有主键,更新会覆盖已有的行,删除会移除对应的行。
数据投递语义为 at-least-once:发生故障后,部分变更可能会被再次写入,而主键表保证重复写入不会产生影响。
Pipeline 不会同步 PostgreSQL 的表结构变更。参见变更表结构。
准备工作
您需要:
- PostgreSQL 10 或更高版本。本文使用 PostgreSQL 内置的
pgoutput逻辑解码插件,无需安装服务端扩展。 - 一个 StarRocks 集群。
- 一个 Flink 1.20 集群,以及运行它所需的 Java 11 或 17。
- 每个 Flink TaskManager 都能通过网络访问 PostgreSQL 的端口(默认
5432),并按照连接 StarRocks 中的说明访问 StarRocks。 - 要同步的每张 PostgreSQL 表都有主键。StarRocks 会使用相同的主键创建目标表。
第一步:准备 PostgreSQL
请以 PostgreSQL 超级用户或表的所有者身份执行以下步骤。示例同步的是 shop 数据库中 public schema 下的表。
-
开启逻辑复制。在
postgresql.conf中将wal_level设置为logical:wal_level = logical
# 每个运行中的 pipeline 使用一个复制槽和一个 WAL sender。
max_replication_slots = 4
max_wal_senders = 4修改
wal_level后需要重启 PostgreSQL。对于 Amazon RDS、Cloud SQL 等托管服务,请在实例的参数组中设置对应的参数(例如rds.logical_replication = 1)。重启后检查该设置:
SHOW wal_level;wal_level
-----------
logical -
为 Flink CDC 创建用户。该用户需要
REPLICATION属性以及表的读取权限:CREATE ROLE flink_cdc WITH LOGIN REPLICATION PASSWORD '<password>';
GRANT CONNECT ON DATABASE shop TO flink_cdc;
GRANT USAGE ON SCHEMA public TO flink_cdc;
GRANT SELECT ON ALL TABLES IN SCHEMA public TO flink_cdc;如果
pg_hba.conf限制了复制连接,请允许该用户从 Flink TaskManager 所在主机连接。例如:# TYPE DATABASE USER ADDRESS METHOD
host shop flink_cdc 10.0.0.0/24 scram-sha-256
host replication flink_cdc 10.0.0.0/24 scram-sha-256 -
为要同步的表创建 publication。Publication 必须包含 pipeline 读取的所有表(即第四步中的
tables配置项):CREATE PUBLICATION flink_pub FOR TABLE public.orders, public.customers;提示请自行创建 publication。如果 publication 不存在,connector 会尝试执行
CREATE PUBLICATION ... FOR ALL TABLES,而只有超级用户才能执行该语句。对于其他用户,该语句会失败,Flink 会不断重试,直到报告的错误变成max_wal_senders,而不是缺少 publication。 -
将每张表的 replica identity 设置为
FULL,使 WAL 为每条 UPDATE 和 DELETE 记录完整的旧行:ALTER TABLE public.orders REPLICA IDENTITY FULL;
ALTER TABLE public.customers REPLICA IDENTITY FULL;请在启动 pipeline 之前完成此设置。使用默认的 replica identity 时,第一条 UPDATE 或 DELETE 会导致 Flink 作业因
NullPointerException失败,并且作业每次重启后都会在这条变更上再次失败。
第二步:准备 StarRocks
创建目标数据库以及 pipeline 使用的用户:
CREATE DATABASE shop;
CREATE USER flink_sink IDENTIFIED BY '<password>';
GRANT CREATE TABLE ON DATABASE shop TO USER flink_sink;
GRANT SELECT, INSERT, UPDATE, DELETE, ALTER ON ALL TABLES IN DATABASE shop TO USER flink_sink;