文档首页/ 数据仓库服务 DWS/ 最佳实践/ 数据迁移/ 使用Flink实时消费DWS Binlog数据最佳实践
更新时间:2026-09-07 GMT+08:00
分享

使用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 补丁版本,以获取最新的性能优化。

    可通过以下命令检查版本号:

    1
    SELECT version();
    
  • 确认GUC参数enable_hstore_binlog_table(该参数控制是否可以创建Binlog表)已开启,9.1.0.200以上实时数仓集群默认已打开,如未打开请联系技术支持。
  • Binlog通过设置enable_binlogenable_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设计为与原表对应的一张辅助表,与原表数据变更保持原子性。

图1 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'
);
表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数据示例

表3 DWS Binlog相关函数

函数名称

描述

pgxc_get_binlog_sync_point(rel_name regclass, slot_name text, checkpoint bool, node_id int)

用于从pg_binlog_slots系统表上获取槽位对应的同步点信息。

  • rel_name: 表名。
  • slot_name: 槽位名。
  • checkpoint: 是否属于checkpoint。
  • node_id: 目标节点的id,为0表示所有节点。

pgxc_get_binlog_changes(rel_name regclass, node_id int, start_csn bigint, end_csn bigint)

获取目标表在指定DN上指定同步点区间的增量数据,(node_id给0表示指定所有DN)。

  • rel_name: 表名
  • node_id: 节点ID
  • start_csn: 起始同步点
  • end_csn:终止同步点

pgxc_register_binlog_sync_point (rel_name regclass, slot_name text, node_id int, end_csn bigint, checkpoint bool, xmin bigint)

用于登记同步点以及checkPoint的同步点位。

  • rel_name:表名
  • slot_name:槽位名
  • node_id:节点ID
  • end_csn:要登记的最新同步点
  • checkpoint:是否属于checkpoint的同步点
  • xmin:同步点对应的xmin

pgxc_consumed_binlog_records (rel_name regclass, node_id int)

用于判断该表对应的binlog是否都被消费完,判断标准:max(binlog辅助表中的csn) <= min(pg_binlog_slots系统表中的syncPoint)。

  • rel_name: 表名
  • node_id: 目标节点的ID, 0表示所有节点

Binlog消费模式

全增量一体化消费

DWS Binlog默认进行全增量一体化消费,Flink任务启动时自动完成“全量快照 -> 增量切换”的无缝衔接。

  • 优势:使用简单,无需额外操作;缺少位点时会进行全量同步,无需担心丢数问题。
  • 劣势:对于数据量极大的表,全量同步的速度较慢,耗时可能在数小时级别,可通过使用pgxc_register_full_sync_point函数提前注册可用点位减少影响(但这样做的话,需要手动补数到当前时间)。
  • 机制:Flink Source 在启动时自动拉取全量数据,并记录当前 Binlog 位点,随后自动切换至增量消费模式,如果已经存在可用位点,就使用该位点进行增量同步。

示例

  1. 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

  2. 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'
    );
    

  3. Flink SQL参考。

    1
    INSERT INTO sink SELECT * FROM test_binlog_source;
    

指定时间戳消费

如果Binlog表使用enable_binlog_timestamp开启Binlog,那么就支持指定时间戳同步数据,通过在Flink SQL中指定binlogStartTime参数可以跳过全量同步阶段并从指定时间戳开始增量同步。

使用方法:Flink SQL需要额外添加表级参数binlogStartTime指定Binlog同步开始的时间戳。

  • 优势:只要能够记录上一次同步的数据时间,就能通过指定该时间实现不丢数的快速取数。
  • 劣势:每次任务重启都要修改该参数值,否则可能会重复入数。
  • 机制:Flink Source 在启动时使用binlogStartTime确定增量消费起点,随后切换至增量消费模式进行消费。

示例

  1. 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);
    

  2. 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时间戳
    );
    

  3. Flink SQL参考:

    1
    INSERT INTO sink SELECT * FROM test_binlog_source;
    

自动跳过全量的增量消费

通过在Flink SQL中开启binlogSkipFullSync参数可以跳过全量同步阶段并从当前时间开始增量同步。

使用方法:Flink SQL需要添加表级参数binlogSkipFullSync且设置为true。

  • 优势:当同步任务重启时可以快速恢复任务。
  • 劣势:故障点和启动时间之间可能存在入数,需要手动补数。
  • 机制:Flink Source 在启动时直接使用当前时间作为增量消费起点,随后切换至增量消费模式进行消费。

示例

  1. 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);
    

  2. 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功能跳过全量同步
    );
    

  3. Flink SQL参考。

    1
    INSERT INTO sink SELECT * FROM test_binlog_source;
    

MRS Flink实时消费DWS Binlog数据示例

下面以MRS的Flink 1.20版本为例,演示实时消费DWS Binlog数据的过程,将DWS数据库中orders_source表数据实时同步到下游表orders_sink

  1. 创建DWS实时数仓集群。为确保网络互通,请将DWS集群跟后面的MRS集群创建在同一个区域可用区虚拟私有云下。

    创建集群界面,选择存算一体1:4云盘规格,这部分规格创建的DWS集群具备实时数仓能力。

  2. 创建带Flink组件的MRS集群,为确保网络互通,请将MRS集群跟前面的DWS集群创建在同一个区域可用区虚拟私有云下。

    本例选择Flink 1.20,对应MRS集群版本可以选择MRS 3.6.0-LTS.1(请以MRS界面实际显示信息为准),MRS各组件对应关系参见MRS组件版本一览表,创建MRS集群过程参见自定义购买MRS集群

  3. 使用客户端连接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秒)
    );
    

  4. 继续在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
    );
    

  5. 对于开启了Binlog的表,并不会立刻记录入库操作的Binlog,需要再给同步任务注册同步点后,才会开始记录Binlog。

    开启Flink同步Binlog任务后,会自动循环进行获取同步点、获取增量数据、注册同步点操作。

    以下步骤演示,在无Flink消费场景下,通过调用DWS自带的Binlog函数模拟获取同步点查询Binlog改动、注册同步点查询Binlog是否已被下游消费的过程。

    获取同步点。先模拟Flink调用系统函数pgxc_get_binlog_sync_point获取同步点,参数分别表示表名、槽位名、是否checkPoint点位,目标DN(为0表示所有DN)。

    1
    SELECT * FROM pgxc_get_binlog_sync_point('orders_source', 'slot1', false, 0);
    

    返回字段说明:

    • node_name: 节点名
    • node_id: 节点ID
    • last_sync_point: 上次同步点
    • latest_sync_point: 当前最新同步点
    • xmin: 同步点对应的xmin

  6. 源表进行增删改操作,产生增量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;
    

  7. 查询Binlog改动。模拟Flink调用系统函数pgxc_get_binlog_changes查询指定CSN区间的Binlog,参数分别表示表名,目标DN(为0表示所有DN),起始CSN点位, 终止CSN点位。

    1
    SELECT * FROM pgxc_get_binlog_changes('orders_source', 0, 0 , 9999999999);
    

    返回字段说明:

    • gs_binlog_sync_point:同步点
    • gs_binlog_event_sequence:用于表示同一事务内的先后顺序
    • gs_binlog_event_typeBinlog类型
    • 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'记录,分别表示更新前的值以及更新后的值。

  8. 注册同步点。模拟Flink调用系统函数pgxc_register_binlog_sync_point注册同步点,参数分别表示表名,槽位名,注册的点位,是否属于checkPoint, 点位对应的xmin (获取同步点时会提供)。

    1
    SELECT pgxc_register_binlog_sync_point('orders_source', 'slot1', 0, 9999999999, false, 100);
    

    返回值: 登记成功的节点数量。

  9. 查询Binlog是否已被下游消费。查询表上Binlog是否被全部消费,返回1表示已经被下游槽位全部消费。

    1
    SELECT * FROM pgxc_consumed_binlog_records('orders_source',0);
    

    这里返回0,表示未被下游槽位消费。下面将继续演示,通过MRS Flink消费Binlog数据的过程

  10. 创建具有FlinkServer相关权限的用户。

    1. 登录MRS控制台,单击已创建的MRS集群名称。
    2. 右侧“运维管理”的“IAM用户同步”,单击“同步”。如果已同步,则跳过该步骤。
    3. 单击“前往Manager”,在弹出的窗口中选择“EIP访问”并配置弹性IP信息,单击“确定”,进入Manager登录页面。如果没有可用EIP,需先跳转到EIP页面创建弹性IP。
    1. Manager登录页面,输入默认用户名“admin”及创建MRS集群时设置的密码,单击“登录”进入Manager页面。
    2. 选择“系统 > 权限 > 角色”。
    3. 单击“添加角色”,在“角色名称”输入角色名字flink01,在“配置资源权限”的表格中选择“待操作集群名称 > Flink”,勾选“FlinkServer管理操作权限”,单击“确定”,返回角色管理。

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

    5. 单击“确定”,成功创建flinktest用户。

  11. 创建Flink作业。

    1. 回到MRS控制台,单击已创建的MRS集群名称。
    2. 单击“前往Manager”,注意如果前期已用admin登录,需要登出,重新以10新创建的flink用户flinktest进行登录。

      如果是初次登录需要修改密码。

    3. 登录Manager页面后,选择“集群 > 服务 > Flink”进入Flink组件页面。
    4. 单击Flink Web UI的路径,访问Flink Web UI。

    5. 进入Flink Web UI后,选择“作业管理 > 新建作业”。
    6. “类型”选择“Flink SQL”,名称可自定义为dws_binlog,作业类型选择“流作业”,单击“确定”。

    7. 进入作业编辑页面,复制以下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;
      

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

  12. 回到连接DWS的客户端界面,查询目标表,显示数据同步成功。

    1
    SELECT * FROM orders_sink;
    

  13. 验证数据是否实时同步。

    1. 插入表orders_source两条新的数据。
      1
      2
      INSERT INTO orders_source VALUES(300, 1, 1);
      INSERT INTO orders_source VALUES(400, 1, 1);
      
    2. 查看下游的目标表是否实时同步。

      结果显示,数据成功同步。

      1
      SELECT * 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 流程
    1. 确认 DWS 实例恢复正常。
    2. 暂停所有下游 Flink 任务。
    3. 触发 Flink 任务全量重启(scan.startup.mode='initial')。
    4. 验证数据一致性后恢复业务。
  • 备份:Binlog 本身不支持独立备份,依赖 DWS 实例级备份与恢复。

安全与权限管控

  • 最小权限:为 Flink 消费账号仅授予 SELECT 及 REPLICATION 权限,禁止 INSERT/UPDATE/DELETE。
  • 透明加密:若 DWS 开启 TDE,Binlog 消费过程自动加解密,对 Flink 透明,无需额外配置。
  • 网络隔离:建议通过 VPC Endpoint 或专线连接 DWS,避免公网暴露。

常见问题与排查指南

表5 常见问题

问题现象

可能原因

排查步骤

消费延迟持续增大

Flink 反压 / DN 负载过高

  1. 检查 Flink Backpressure 面板。
  2. 检查 DWS DN CPU/IO 监控。
  3. 考虑增加 Flink 并发或优化 SQL。

Slot 位点不更新

Flink 任务挂起 / 网络中断

  1. 检查 Flink TaskManager 日志。
  2. 测试 DWS 网络连通性。
  3. 检查 pgxc_get_binlog_consume_progress。

数据丢失/重复

Checkpoint 失败 / 位点重置

  1. 检查 Flink Checkpoint 状态。
  2. 确认是否执行了 DDL 操作。
  3. 核对 ignoreUpdateBefore 配置。

启动报 “Slot not found”

Slot 被清理 / 名称错误

  1. 核对 binlogSlotName 拼写。
  2. 确认该slot是否因为超出TTL被清理。
  3. 重新触发全量同步。

相关文档