使用Flink实时消费DWS Binlog数据最佳实践
当用户需要捕获数据库事件,将增量数据导出至Flink等第三方组件,并协同完成数据加工等任务时,可通过DWS实时数仓中HStore表的Binlog功能,消费Binlog数据来实现上下游的数据同步,提高数据加工的效率。
传统的数据库如MySQL,支持通过Binlog来记录数据库中所有数据的变化,MySQL的Binlog主要用于数据恢复与主从复制,而DWS实时数仓的Binlog功能一般只用于实时场景下的数据同步。同时DWS实时数仓Binlog并不会记录DDL操作,只记录INSERT/DELETE/UPDATE(UPSERT)等DML操作。
DWS的Binlog的优点如下:
- 表级按需开关:按需给指定表打开或关闭Binlog功能,更为灵活。
- 全增量一体消费:支持Flink任务启动后,先全量同步源端数据,再实时消费源端增量数据。
- 支持消费即清理:对于空间敏感且只关注实时同步与加工的客户,支持消费后即开始异步清理增量,有效减少空间使用。
- 快速构建:利用Flink强大的实时处理能力和DWS的Binlog能力,可以快速构建实时数仓,且无需维护其他组件(如Kafka),整体架构分层清晰,数据可以高效流动,并且整体任务链路都可以通过Flink SQL来驱动,便于业务人员理解和使用。
本实践介绍DWS Binlog基本功能,并以MRS的Flink 1.20版本为例,演示如何实时消费DWS Binlog数据,将DWS数据库中orders_source表的数据实时同步到下游表orders_sink中,详见MRS Flink实时消费DWS Binlog数据示例。
基本概念
- Binlog:即二进制日志(Binary Log),这个概念常见于MySQL数据库,是用来记录数据库所有可能引起数据变化事件的日志,比如表结构变更(例如CREATE、ALTER TABLE)以及表数据修改(INSERT、UPDATE、DELETE)事件。
- HStore表:HStore是DWS实时数仓提供的一种高性能列式存储表引擎,支持高效的UPSERT(插入或更新)操作。开启Binlog功能的前提是表必须为HStore表。
- CDC:CDC(Change Data Capture,变更数据捕获)是一种识别并响应数据变更的设计模式。DWS的Binlog功能与Flink CDC结合后,Source端会根据操作类型自动为每行数据设置准确的Flink RowKind类型(INSERT、DELETE、UPDATE_BEFORE、UPDATE_AFTER),这样就能镜像同步表的数据,类似MySQL和PostgreSQL的CDC功能。
- 复制槽(Slot):复制槽(Slot)是DWS数据库DN节点上持久保存的消费进度标记,专门用于CDC、逻辑解码、增量同步,记录外部消费者(Flink‑CDC)读到 WAL日志哪一个位点(LSN/CSN)。在DWS的Binlog函数中,使用binlogSlotName表示槽位名称,槽位是DWS Binlog消费中的一个逻辑标识,用于区分不同的消费任务。由于可能存在多个Flink任务同时消费同一张表的Binlog信息,所以该场景需要保证每个任务的binlogSlotName不同。
- 同步点(Sync Point):同步点是Binlog消费中的位置标记,表示消费任务已经读取到哪个位置。
约束限制
- 为确保功能的完整性与稳定性,使用DWS Binlog功能需满足以下版本要求:
- 最低版本:DWS 9.1.0.200 及以上版本。
- 推荐版本:建议升级至当前最新的 9.1.x 补丁版本,以获取最新的性能优化。
可通过以下命令检查版本号:
1SELECT version();
- 确认GUC参数enable_hstore_binlog_table(该参数控制是否可以创建Binlog表)已开启,9.1.0.200以上实时数仓集群默认已打开,如未打开请联系技术支持。
- Binlog通过设置enable_binlog或enable_binlog_timestamp表级参数开启,且开启Binlog的表必须满足以下条件:
- 必须是HStore或HStore_opt表。
- 必须包含主键。
- 必须使用Hash分布模式。
- Binlog只记录INSERT/UPDATE/UPSERT/DELETE等DML操作,不记录DDL操作,且进行部分DDL操作(ADD COLUMN增加列、DROP COLUMN删除列、SET TYPE修改列、TRUNCATE清空表数据)会将增量数据和数据同步位点全部清除。若必须执行DDL,需要停止上游数据写入,待存量数据全部写入后再执行DDL,以防止数据丢失。
- Binlog只用于实时数据同步或增量数据加工,严禁用于审计日志或历史数据归档等功能。
- Binlog表名不要使用特殊字符,例如:点(.)、引号("" 或 '')等。
- Binlog表当前不支持以下操作:INSERT OVERWRITE、修改分布列、给临时表开启Binlog、EXCHANGE/MERGE/SPLIT PARTITION。
- Binlog表在扩缩容或VACUUM FULL操作期间会等待存量Binlog记录的消费,全部消费完成才会进行对应步骤,默认等待时间为一小时,如果等待超时或者出错就会退出操作过程,视为操作失败。即使VACUUM FULL的对象只有一个分区,因为需要等待Binlog消费,也会对主表加锁,阻塞整个表的插入、更新或删除。
- Binlog表在备份恢复期间会被当做普通HStore表进行备份,恢复后辅助表的增量数据及数据同步位点会清空,需要重新同步。
架构规划与表设计原则
DWS Binlog是一种内嵌于存储引擎的变更数据捕获(CDC)机制。与传统数据库依赖外部日志文件(如 MySQL Binlog、PostgreSQL WAL)不同,DWS将Binlog设计为与原表对应的一张辅助表,与原表数据变更保持原子性。
Flink表使用介绍
使用Flink SQL的表示例如下,参数解释参见表1,了解更多请参见Flink实时消费Binlog。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 | CREATE TABLE dws_source ( user_id BIGINT, event_time TIMESTAMP(3), amount DECIMAL(10, 2), PRIMARY KEY (user_id) NOT ENFORCED ) WITH ( 'connector' = 'dws', 'url' = 'jdbc:gaussdb://192.168.0.1:8000/gaussdb', 'binlog' = 'true', 'username' = 'flink_user', 'password' = '***', 'tableName' = 'test_binlog', 'binlogSlotName' = 'binlog_job_1', 'connectionSize' = '1' ); |
| 参数 | 说明 | 默认值 |
|---|---|---|
| binlog | 是否读取Binlog信息。 | false |
| binlogSlotName | 槽位信息,可以理解一个标识。由于可能存在多个Flink任务同时消费同一张表的Binlog信息,所以该场景需要保证每个任务的binlogSlotName不同。 | Flink映射表的表名 |
| fullSyncBinlogReadTimeout | 全量消费Binlog数据超时时间,单位毫秒。 | 1800000 |
| binlogReadTimeout | 增量消费Binlog数据超时时间,单位毫秒。 | 600000 |
| binlogSleepTime | 实时消费不到Binlog数据休眠时间,单位毫秒。 | 500 |
| binlogMaxRetryTimes | 消费Binlog数据出错后的重试次数。 | 1 |
| connectionPoolSize | JDBC连接池连接大小。 | 5 |
| binlogSyncPointSize | 增量读取Binlog同步点区间的大小。 | 5000 |
| binlogQueueSize | Binlog增量读取时队列长度。 | 0(无限制) |
Flink任务并发推荐
- 连接数(Flink任务并发度 * Flink表级参数connectionPoolSize),建议参考以下数据进行配置。
表2 连接数 表流量(条/s)
binlogSyncPointSize
连接数(并发数)
20000+
10w
DN数
6000~20000
10w
DN数/2
3000~6000
10w
DN数/4
1000~3000
100w
DN数/8
0~1000
100w
DN数/16
- 连接数越高性能越好,资源足够的场景可以将连接数上调。
- 集群Binlog并发推荐所有Binlog任务并发总和不超过CN的max_connections参数值的80%(与集群业务模型有关,普通业务占比高时可根据实际业务比例适当下调)。
DWS Binlog相关函数
DWS的Binlog表相关函数如表3所示,使用示例参见MRS Flink实时消费DWS Binlog数据示例。
| 函数名称 | 描述 |
|---|---|
| pgxc_get_binlog_sync_point(rel_name regclass, slot_name text, checkpoint bool, node_id int) | 用于从pg_binlog_slots系统表上获取槽位对应的同步点信息。
|
| pgxc_get_binlog_changes(rel_name regclass, node_id int, start_csn bigint, end_csn bigint) | 获取目标表在指定DN上指定同步点区间的增量数据,(node_id给0表示指定所有DN)。
|
| 用于登记同步点以及checkPoint的同步点位。
| |
| pgxc_consumed_binlog_records (rel_name regclass, node_id int) | 用于判断该表对应的binlog是否都被消费完,判断标准:max(binlog辅助表中的csn) <= min(pg_binlog_slots系统表中的syncPoint)。
|
Binlog消费模式
全增量一体化消费
DWS Binlog默认进行全增量一体化消费,Flink任务启动时自动完成“全量快照 -> 增量切换”的无缝衔接。
- 优势:使用简单,无需额外操作;缺少位点时会进行全量同步,无需担心丢数问题。
- 劣势:对于数据量极大的表,全量同步的速度较慢,耗时可能在数小时级别,可通过使用pgxc_register_full_sync_point函数提前注册可用点位减少影响(但这样做的话,需要手动补数到当前时间)。
- 机制:Flink Source 在启动时自动拉取全量数据,并记录当前 Binlog 位点,随后自动切换至增量消费模式,如果已经存在可用位点,就使用该位点进行增量同步。
示例:
- DWS建表。关键参数参见表4
1 2 3 4 5 6 7 8 9 10
CREATE TABLE hstore_binlog_source ( c1 INT PRIMARY KEY, -- 必须有主键 c2 INT, c3 INT ) WITH ( ORIENTATION = COLUMN, enable_hstore_opt=true, -- 必须是HStore/HStore-opt表 enable_binlog=on, -- 开启Binlog binlog_ttl = 259200 ) DISTRIBUTE BY HASH(c1); -- 必须使用HASH分布
表4 Binlog表关键参数说明 关键参数
描述
取值
enable_hstore_opt
是否创建Hstore_opt表,了解更多Hstore_opt参见实时数仓及HStore表使用最佳实践。
true
enable_binlog
用于控制打开Binlog功能,仅对HStore以及HStore_opt表有效,且必须有主键。
on
binlog_ttl
可选参数,当不设置时将使用默认值86400, 单位为秒,当同步任务注册的同步点超过TTL没有进行增量同步时,该同步点位将被清理。最老的同步点位之前的Binlog(即被所有任务消费了的Binlog)会被异步清理来回收空间。
259200
DISTRIBUTE BY
指定表数据如何在节点之间的分布或者复制方式,Binlog表必须使用HASH。
取值范围:
- REPLICATION:表的每一行存储在所有数据节点(DN)中,即每个数据节点都有完整的表数据。
- ROUNDROBIN:表的每一行被轮询发送给各个DN,因此数据会被均匀地分布在各个DN中。(ROUNDROBIN仅8.1.2及以上版本支持)
- HASH (column_name ) :对指定的列进行Hash,通过映射,把数据分布到指定DN。
HASH
- Flink对应建表。
1 2 3 4 5 6 7 8 9 10 11 12 13 14
CREATE TABLE test_binlog_source ( c1 int, c2 int, c3 int, primary key(c1) NOT ENFORCED -- 对应数据表创建主键 ) with ( 'connector' = 'dws', 'url' = 'jdbc:gaussdb://ip:port/database', 'binlog' = 'true', 'tableName' = 'test_binlog_source', 'binlogSlotName' = 'slot', 'username'='xxx', 'password'='xxx' );
- Flink SQL参考。
1INSERT INTO sink SELECT * FROM test_binlog_source;
指定时间戳消费
如果Binlog表使用enable_binlog_timestamp开启Binlog,那么就支持指定时间戳同步数据,通过在Flink SQL中指定binlogStartTime参数可以跳过全量同步阶段并从指定时间戳开始增量同步。
使用方法:Flink SQL需要额外添加表级参数binlogStartTime指定Binlog同步开始的时间戳。
- 优势:只要能够记录上一次同步的数据时间,就能通过指定该时间实现不丢数的快速取数。
- 劣势:每次任务重启都要修改该参数值,否则可能会重复入数。
- 机制:Flink Source 在启动时使用binlogStartTime确定增量消费起点,随后切换至增量消费模式进行消费。
示例:
- DWS建表。
1 2 3 4 5 6 7 8 9 10 11
CREATE TABLE hstore_binlog_source ( c1 INT PRIMARY KEY, c2 INT, c3 INT ) WITH ( ORIENTATION = COLUMN, enable_hstore_opt=true, -- 需要使用enable_binlog_timestamp开启Binlog enable_binlog_timestamp=on, binlog_ttl = 259200 ) DISTRIBUTE BY HASH(c1);
- Flink对应建表。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
CREATE TABLE test_binlog_source ( c1 int, c2 int, c3 int, primary key(c1) NOT ENFORCED ) with ( 'connector' = 'dws', 'url' = 'jdbc:gaussdb://ip:port/database', 'binlog' = 'true', 'tableName' = 'test_binlog_source', 'binlogSlotName' = 'slot', 'username'='xxx', 'password'='xxx', 'binlogStartTime'='0' -- 指定你要开始同步的时间点,unix时间戳 );
- Flink SQL参考:
1INSERT INTO sink SELECT * FROM test_binlog_source;
自动跳过全量的增量消费
通过在Flink SQL中开启binlogSkipFullSync参数可以跳过全量同步阶段并从当前时间开始增量同步。
使用方法:Flink SQL需要添加表级参数binlogSkipFullSync且设置为true。
- 优势:当同步任务重启时可以快速恢复任务。
- 劣势:故障点和启动时间之间可能存在入数,需要手动补数。
- 机制:Flink Source 在启动时直接使用当前时间作为增量消费起点,随后切换至增量消费模式进行消费。
示例:
- DWS建表。
1 2 3 4 5 6 7 8 9 10
CREATE TABLE hstore_binlog_source ( c1 INT PRIMARY KEY, c2 INT, c3 INT ) WITH ( ORIENTATION = COLUMN, enable_hstore_opt=true, enable_binlog =on, binlog_ttl = 259200 ) DISTRIBUTE BY HASH(c1);
- Flink对应建表。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
CREATE TABLE test_binlog_source ( c1 int, c2 int, c3 int, primary key(c1) NOT ENFORCED ) with ( 'connector' = 'dws', 'url' = 'jdbc:gaussdb://ip:port/database', 'binlog' = 'true', 'tableName' = 'test_binlog_source', 'binlogSlotName' = 'slot', 'username'='xxx', 'password'='xxx', 'binlogSkipFullSync'='true' -- 开启skipFullSync功能跳过全量同步 );
- Flink SQL参考。
1INSERT INTO sink SELECT * FROM test_binlog_source;
MRS Flink实时消费DWS Binlog数据示例
下面以MRS的Flink 1.20版本为例,演示实时消费DWS Binlog数据的过程,将DWS数据库中orders_source表数据实时同步到下游表orders_sink。
- 创建DWS实时数仓集群。为确保网络互通,请将DWS集群跟后面的MRS集群创建在同一个区域、可用区及虚拟私有云下。
在创建集群界面,选择存算一体1:4云盘规格,这部分规格创建的DWS集群具备实时数仓能力。
- 创建带Flink组件的MRS集群,为确保网络互通,请将MRS集群跟前面的DWS集群创建在同一个区域、可用区及虚拟私有云下。
本例选择Flink 1.20,对应MRS集群版本可以选择MRS 3.6.0-LTS.1(请以MRS界面实际显示信息为准),MRS各组件对应关系参见MRS组件版本一览表,创建MRS集群过程参见自定义购买MRS集群。
- 使用客户端连接DWS集群,创建源表,并开启Binlog。
1 2 3 4 5 6 7 8 9 10 11
-- 创建一张电商订单表作为Binlog源表 CREATE TABLE orders_source ( c1 INT PRIMARY KEY, -- 订单ID,主键(Binlog表必须存在主键约束) c2 INT, -- 用户ID c3 INT -- 订单金额(以分为单位,避免浮点精度问题) ) WITH ( ORIENTATION = COLUMN, enable_hstore_opt = true, -- 创建Hstore_opt表。 enable_binlog = on, -- 开启普通Binlog功能,记录表上所有DML变更 binlog_ttl = 259200 -- Binlog记录保留时间为3天(259200秒) );
- 继续在DWS数据库中创建目标表,结构与源表保持一致,目标表不需要开启Binlog。
1 2 3 4 5 6 7 8
CREATE TABLE orders_sink ( c1 INT PRIMARY KEY, -- 订单ID,主键(与源表一致) c2 INT, -- 用户ID c3 INT -- 订单金额 ) WITH ( ORIENTATION = COLUMN, enable_hstore_opt = true -- 启用HStore优化版引擎,支持高效UPSERT );
- 对于开启了Binlog的表,并不会立刻记录入库操作的Binlog,需要再给同步任务注册同步点后,才会开始记录Binlog。
开启Flink同步Binlog任务后,会自动循环进行获取同步点、获取增量数据、注册同步点操作。
以下步骤演示,在无Flink消费场景下,通过调用DWS自带的Binlog函数模拟获取同步点、查询Binlog改动、注册同步点、查询Binlog是否已被下游消费的过程。
获取同步点。先模拟Flink调用系统函数pgxc_get_binlog_sync_point获取同步点,参数分别表示表名、槽位名、是否checkPoint点位,目标DN(为0表示所有DN)。
1SELECT * FROM pgxc_get_binlog_sync_point('orders_source', 'slot1', false, 0);

返回字段说明:
- node_name: 节点名。
- node_id: 节点ID。
- last_sync_point: 上次同步点。
- latest_sync_point: 当前最新同步点。
- xmin: 同步点对应的xmin。
- 源表进行增删改操作,产生增量Binlog数据。
1 2 3 4
INSERT INTO orders_source VALUES(100, 1, 1); DELETE FROM orders_source WHERE c1 = 100; INSERT INTO orders_source VALUES(200, 1, 1); UPDATE orders_source SET c2 = 2 WHERE c1 = 200;
- 查询Binlog改动。模拟Flink调用系统函数pgxc_get_binlog_changes查询指定CSN区间的Binlog,参数分别表示表名,目标DN(为0表示所有DN),起始CSN点位, 终止CSN点位。
1SELECT * FROM pgxc_get_binlog_changes('orders_source', 0, 0 , 9999999999);

返回字段说明:
- gs_binlog_sync_point:同步点。
- gs_binlog_event_sequence:用于表示同一事务内的先后顺序。
- gs_binlog_event_type:Binlog类型。
- gs_binlog_timestamp_us:Binlog记录的时间戳,对于enable_binlog_timestamp为false的Binlog表,该列返回空,本例的表定义中未指定enable_binlog_timestamp,所以为空。
- ...value columns...: 目标表上各个用户字段的数据。
从以上回显信息,可以看到两次INSERT操作产生了两个gs_binlog_event_type是'I'的记录,DELETE操作产生了type是'd'的记录,UPDATE产生了一行BeforeUpdate的'B'记录以及一条AfterUpdate的'U'记录,分别表示更新前的值以及更新后的值。
- 注册同步点。模拟Flink调用系统函数pgxc_register_binlog_sync_point注册同步点,参数分别表示表名,槽位名,注册的点位,是否属于checkPoint, 点位对应的xmin (获取同步点时会提供)。
1SELECT pgxc_register_binlog_sync_point('orders_source', 'slot1', 0, 9999999999, false, 100);

返回值: 登记成功的节点数量。
- 查询Binlog是否已被下游消费。查询表上Binlog是否被全部消费,返回1表示已经被下游槽位全部消费。
1SELECT * FROM pgxc_consumed_binlog_records('orders_source',0);

这里返回0,表示未被下游槽位消费。下面将继续演示,通过MRS Flink消费Binlog数据的过程。
- 创建具有FlinkServer相关权限的用户。
- 登录MRS控制台,单击已创建的MRS集群名称。
- 右侧“运维管理”的“IAM用户同步”,单击“同步”。如果已同步,则跳过该步骤。
- 单击“前往Manager”,在弹出的窗口中选择“EIP访问”并配置弹性IP信息,单击“确定”,进入Manager登录页面。如果没有可用EIP,需先跳转到EIP页面创建弹性IP。
- Manager登录页面,输入默认用户名“admin”及创建MRS集群时设置的密码,单击“登录”进入Manager页面。
- 选择“系统 > 权限 > 角色”。
- 单击“添加角色”,在“角色名称”输入角色名字flink01,在“配置资源权限”的表格中选择“待操作集群名称 > Flink”,勾选“FlinkServer管理操作权限”,单击“确定”,返回角色管理。

- 回到Manager页面,选择“系统 > 权限 > 用户”,单击“添加用户”,按如下配置。
- 用户名输入flinktest。
- 用户类型选择“人机”。
- 输入密码。
- 用户组:关联supergroup。
- 角色:绑定新创建的角色flink01和System_administrator。

- 单击“确定”,成功创建flinktest用户。
- 创建Flink作业。
- 回到MRS控制台,单击已创建的MRS集群名称。
- 单击“前往Manager”,注意如果前期已用admin登录,需要登出,重新以10新创建的flink用户flinktest进行登录。
如果是初次登录需要修改密码。
- 登录Manager页面后,选择“集群 > 服务 > Flink”进入Flink组件页面。
- 单击Flink Web UI的路径,访问Flink Web UI。
- 进入Flink Web UI后,选择“作业管理 > 新建作业”。
- “类型”选择“Flink SQL”,名称可自定义为dws_binlog,作业类型选择“流作业”,单击“确定”。
- 进入作业编辑页面,复制以下SQL到作业编辑框中,注意替换DWS集群IP、dbadmin的密码,右侧勾选“开启作业定时调优”,其他默认即可,单击“保存”,再单击“提交”。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39
-- 在Flink SQL中创建源表映射 CREATE TABLE flink_orders_source ( c1 INT, -- 订单ID c2 INT, -- 用户ID c3 INT, -- 订单金额 PRIMARY KEY (c1) NOT ENFORCED -- 声明主键,NOT ENFORCED表示Flink不强制约束,由DWS保证 ) WITH ( 'connector' = 'dws', -- 使用DWS连接器 'url' = 'jdbc:gaussdb://<DWS集群IP>:8000/gaussdb', -- DWS JDBC连接地址 'binlog' = 'true', -- 开启Binlog消费模式(必填) 'tableName' = 'orders_source', -- DWS中的源表名(必填) 'binlogSlotName' = 'flink_slot_1', -- 槽位名,多任务消费同一表时需保证唯一 'executeDirectOn' = 'true', -- Binlog读取优化,全量同步时避免CN结果集下盘,增量同步提升3-4倍性能(仅9.1.0.200+) 'username' = 'dbadmin', -- DWS数据库用户名(必填) 'password' = '<密码>' -- DWS数据库密码(必填) ); -- 在Flink SQL中创建目标表映射 CREATE TABLE flink_orders_sink ( c1 INT, -- 订单ID c2 INT, -- 用户ID c3 INT, -- 订单金额 PRIMARY KEY (c1) NOT ENFORCED ) WITH ( 'connector' = 'dws', -- 使用DWS连接器 'url' = 'jdbc:gaussdb://<DWS集群IP>:8000/gaussdb', -- DWS JDBC连接地址 'tableName' = 'orders_sink', -- DWS中的目标表名 'ignoreUpdateBefore' = 'false', -- 是否忽略UPDATE_BEFORE记录 'connectionSize' = '1', -- 连接数设为1,保证DN内数据写入顺序 'username' = 'dbadmin', -- DWS数据库用户名 'password' = '<密码>' -- DWS数据库密码 ); -- 启动同步任务,源表数据同步到目标表。 -- Flink会自动完成以下流程: INSERT INTO flink_orders_sink SELECT * FROM flink_orders_source;

回到作业列表,作业会自动运行,显示“运行中”,表示作业正常运行。

- 回到连接DWS的客户端界面,查询目标表,显示数据同步成功。
1SELECT * FROM orders_sink;

- 验证数据是否实时同步。
稳定性保障与运维监控
监控消费进度
DWS Binlog通过pgxc_get_binlog_consume_progress(rel_name text, node_id int)函数监控Binlog消费进度。
描述:该函数用于获取目标表的消费进度信息,可以实时监控各 Slot 的消费状态,只能对开启Binlog时间戳的表使用。
返回值类型:record
返回值:
- node_name:节点名。
- node_id:节点ID。
- slot_name:槽位名。
- checkpoint:是否属于checkpoint点位。
- latest_consumed_timestamp:已消费的最新Binlog的时间戳。
- latest_timestamp:所有Binlog里最新的时间戳。
- latest_consumed_csn: 已消费的最新Binlog的csn点位。
- latest_csn:所有Binlog里最新的csn点位。
- unconsumed_binlog_count:还未消费的Binlog的数量。
示例:
1 2 3 4 5 6 | SELECT * FROM pgxc_get_binlog_consume_progress('hstore_binlog_source', 0); node_name | node_id | slot_name | checkpoint | latest_consumed_timestamp | latest_timestamp | latest_consumed_csn | latest_csn | unconsumed_binlog_count -----------+-------------+-----------+------------+---------------------------+------------------------+---------------------+------------+------------------------- dn_1 | -1300059100 | slot1 | f | 2025-04-15 15:55:55+08 | 2025-04-15 15:55:55+08 | 417725 | 417725 | 0 dn_1 | -1300059100 | slot1 | t | 2025-04-15 15:55:55+08 | 2025-04-15 15:55:55+08 | 417725 | 417725 | 0 (2 rows) |
建议的告警策略:
- 延迟告警:unconsumed_binlog_count > 10000 或 latest_csn – latest_consumed_csn > 100000。
- 停滞告警:latest_consumed_timestamp 超过 5 分钟未更新。
数据恢复预案
- 实例恢复:DWS 实例故障恢复后,Binlog 增量数据与 Slot 详情会被重置。
- SOP 流程:
- 确认 DWS 实例恢复正常。
- 暂停所有下游 Flink 任务。
- 触发 Flink 任务全量重启(scan.startup.mode='initial')。
- 验证数据一致性后恢复业务。
- 备份:Binlog 本身不支持独立备份,依赖 DWS 实例级备份与恢复。
安全与权限管控
- 最小权限:为 Flink 消费账号仅授予 SELECT 及 REPLICATION 权限,禁止 INSERT/UPDATE/DELETE。
- 透明加密:若 DWS 开启 TDE,Binlog 消费过程自动加解密,对 Flink 透明,无需额外配置。
- 网络隔离:建议通过 VPC Endpoint 或专线连接 DWS,避免公网暴露。
常见问题与排查指南
| 问题现象 | 可能原因 | 排查步骤 |
|---|---|---|
| 消费延迟持续增大 | Flink 反压 / DN 负载过高 |
|
| Slot 位点不更新 | Flink 任务挂起 / 网络中断 |
|
| 数据丢失/重复 | Checkpoint 失败 / 位点重置 |
|
| 启动报 “Slot not found” | Slot 被清理 / 名称错误 |
|


