# 使用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 补丁版本，以获取最新的性能优化。
  
  
  可通过以下命令检查版本号：
  ```
  SELECT 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设计为与原表对应的一张辅助表，与原表数据变更保持原子性。
图1DWS Binlog入库原理   
![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002747303689.png "点击放大")
#### Flink表使用介绍
使用Flink SQL的表示例如下，参数解释参见[表1]，了解更多请参见[Flink实时消费Binlog](https://support.huaweicloud.com/devg-911-dws/dws_04_1212.html)。
```
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数据示例]。
 表3DWS Binlog相关函数 
| 函数名称                                                                                                                                                                                         | 描述                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                         |
|:---|:---|
| [pgxc_get_binlog_sync_point(rel_name regclass, slot_name text, checkpoint bool, node_id int)](https://support.huaweicloud.com/devg-dws/dws_04_1027.html)                                    | 用于从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)](https://support.huaweicloud.com/devg-dws/dws_04_1027.html)                                      | 获取目标表在指定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)](https://support.huaweicloud.com/devg-dws/dws_04_1027.html) | 用于登记同步点以及checkPoint的同步点位。 - rel_name：表名。  - slot_name：槽位名。  - node_id：节点ID。  - end_csn：要登记的最新同步点。  - checkpoint：是否属于checkpoint的同步点。  - xmin：同步点对应的xmin。   |
| [pgxc_consumed_binlog_records (rel_name regclass, node_id int)](https://support.huaweicloud.com/devg-dws/dws_04_1027.html)                                                                  | 用于判断该表对应的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]
   
   ```
   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分布
   ```
    表4Binlog表关键参数说明 
   | 关键参数            | 描述                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                      | 取值     |
   |:---|:---|:---|
   | enable_hstore_opt | 是否创建Hstore_opt表，了解更多Hstore_opt参见[实时数仓及HStore表使用最佳实践](https://support.huaweicloud.com/bestpractice-dws/dws_05_0029.html)。                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                                 | 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对应建表。 
   ```
   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参考。 
   ```
   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建表。 
   ```
   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对应建表。 
   ```
   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参考： 
   ```
   INSERT INTO sink SELECT * FROM test_binlog_source;
   ```
   
   
 
#### 自动跳过全量的增量消费
通过在Flink SQL中开启binlogSkipFullSync参数可以跳过全量同步阶段并从当前时间开始增量同步。
**使用方法**：Flink SQL需要添加表级参数binlogSkipFullSync且设置为true。
- **优势**：当同步任务重启时可以快速恢复任务。
- **劣势**：故障点和启动时间之间可能存在入数，需要手动补数。
- **机制**：Flink Source 在启动时直接使用当前时间作为增量消费起点，随后切换至增量消费模式进行消费。
**示例**：
1. DWS建表。 
   ```
   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对应建表。 
   ```
   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参考。 
   ```
   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集群创建在**同一个区域** 、**可用区** 及**虚拟私有云** 下。
   
   
   在[创建集群](https://support.huaweicloud.com/mgtg-dws/dws_01_0019.html)界面，选择**存算一体1:4云盘**规格，这部分规格创建的DWS集群具备实时数仓能力。
   
   
2. 创建带**Flink** 组件的MRS集群，为**确保网络互通** ，请将MRS集群跟前面的DWS集群创建在**同一个区域** 、**可用区** 及**虚拟私有云** 下。
   
   本例选择**Flink 1.20** ，对应MRS集群版本可以选择**MRS 3.6.0-LTS.1** （请以MRS界面实际显示信息为准），MRS各组件对应关系参见[MRS组件版本一览表](https://support.huaweicloud.com/productdesc-mrs/mrs_08_0005.html)，创建MRS集群过程参见[自定义购买MRS集群](https://support.huaweicloud.com/usermanual-mrs/mrs_01_0513.html)。
   
   
3. 使用客户端连接DWS集群，创建源表，并开启Binlog。 
   ```
   -- 创建一张电商订单表作为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。 
   ```
   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](https://support.huaweicloud.com/devg-dws/dws_04_1027.html)获取同步点，参数分别表示表名、槽位名、是否checkPoint点位，目标DN(为0表示所有DN)。
   ```
   SELECT * FROM pgxc_get_binlog_sync_point('orders_source', 'slot1', false, 0);
   ```
   ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002742073989.png "点击放大")
   返回字段说明：
   - node_name: 节点名。
   
   - node_id: 节点ID。
   
   - last_sync_point: 上次同步点。
   
   - latest_sync_point: 当前最新同步点。
   
   - xmin: 同步点对应的xmin。
   
   
   
   
6. 源表进行增删改操作，产生增量Binlog数据。 
   ```
   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](https://support.huaweicloud.com/devg-dws/dws_04_1027.html)查询指定CSN区间的Binlog，参数分别表示表名，目标DN(为0表示所有DN)，起始CSN点位， 终止CSN点位。
   
   ```
   SELECT * FROM pgxc_get_binlog_changes('orders_source', 0, 0 , 9999999999);
   ```
   ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002742076597.png "点击放大")
   返回字段说明：
   - 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'记录**，分别表示更新前的值以及更新后的值。
   
   
8. **注册同步点** 。模拟Flink调用系统函数[pgxc_register_binlog_sync_point](https://support.huaweicloud.com/devg-dws/dws_04_1027.html)注册同步点，参数分别表示表名，槽位名，注册的点位，是否属于checkPoint, 点位对应的xmin (获取同步点时会提供)。
   
   ```
   SELECT pgxc_register_binlog_sync_point('orders_source', 'slot1', 0, 9999999999, false, 100);
   ```
   ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002712366712.png "点击放大")
   **返回值** **:**登记成功的节点数量。
   
   
9. **查询Binlog是否已被下游消费** 。查询表上Binlog是否被全部消费，返回1表示已经被下游槽位全部消费。
   
   ```
   SELECT * FROM pgxc_consumed_binlog_records('orders_source',0);
   ```
   ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002741969685.png "点击放大")
   这里返回0，表示未被下游槽位消费。**下面将继续演示，通过MRS Flink消费Binlog数据的过程**。
   
   
10. 创建具有FlinkServer相关权限的用户。
    
    1. 登录MRS控制台，单击已创建的MRS集群名称。
    
    2. 右侧"运维管理"的"IAM用户同步"，单击"同步"。如果已同步，则跳过该步骤。
    
    3. 单击"前往Manager"，在弹出的窗口中选择"EIP访问"并配置弹性IP信息，单击"确定"，进入Manager登录页面。如果没有可用EIP，需先跳转到EIP页面创建弹性IP。
    
    
    
    4. Manager登录页面，输入默认用户名"admin"及创建MRS集群时设置的密码，单击"登录"进入Manager页面。
    
    5. 选择"系统 \> 权限 \> 角色"。
    
    6. 单击"添加角色"，在"角色名称"输入角色名字flink01，在"配置资源权限"的表格中选择"*待操作集群名称* \> Flink"，勾选"FlinkServer管理操作权限"，单击"确定"，返回角色管理。
       ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002713202734.png "点击放大")
       
       
    
    7. 回到Manager页面，选择"系统 \> 权限 \> 用户"，单击"添加用户"，按如下配置。
       - 用户名输入flinktest。
       
       - 用户类型选择"人机"。
       
       - 输入密码。
       
       - 用户组：关联**supergroup**。
       
       - 角色：绑定新创建的角色**flink01** 和**System_administrator**。
       
       
       ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002743043521.png "点击放大")
       
       
    
    8. 单击"确定"，成功创建**flinktest**用户。
    
    
    
    
11. 创建Flink作业。 
    1. 回到MRS控制台，单击已创建的MRS集群名称。
    
    2. 单击"前往Manager"，注意如果前期已用admin登录，需要登出，重新以[10]新创建的flink用户**flinktest** 进行登录。
       如果是初次登录需要修改密码。
       
    
    3. 登录Manager页面后，选择"集群 \> 服务 \> Flink"进入Flink组件页面。
    
    4. 单击Flink Web UI的路径，访问Flink Web UI。 ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002743177427.png "点击放大")
       
       
    
    5. 进入Flink Web UI后，选择"作业管理 \> 新建作业"。
    
    6. "类型"选择"Flink SQL"，名称可自定义为dws_binlog，作业类型选择"流作业"，单击"确定"。 ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002744111487.png "点击放大")
       
    
    7. 进入作业编辑页面，复制以下SQL到作业编辑框中，注意替换DWS集群IP、dbadmin的密码，右侧勾选"开启作业定时调优"，其他默认即可，单击"保存"，再单击"提交"。
       ```
       -- 在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;
       ```
       ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002743942605.png "点击放大")
       
       回到作业列表，作业会自动运行，显示"运行中"，表示作业正常运行。
       ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002714504762.png "点击放大")
       
    
    
    
    
12. 回到连接DWS的客户端界面，查询目标表，显示数据同步成功。 
    ```
    SELECT * FROM orders_sink;
    ```
    ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002743945929.png)
    
    
    
13. 验证数据是否实时同步。 
    1. 插入表orders_source两条新的数据。
       ```
       INSERT INTO orders_source VALUES(300, 1, 1);
       INSERT INTO orders_source VALUES(400, 1, 1);
       ```
       
    
    2. 查看下游的目标表是否实时同步。 结果显示，数据成功同步。
       ```
       SELECT * FROM orders_sink;
       ```
       ![](https://support.huaweicloud.com/bestpractice-dws/figure/zh-cn_image_0000002714347910.png)
       
       
    
    
    
    
 
#### 稳定性保障与运维监控
#### 监控消费进度
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的数量。
**示例**：
```
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. 重新触发全量同步。                          |
   
