Flink CDC 实时捕获数据库变更

FreeGuideOnline 最新 2026-07-10

sql -- 检查 binlog 是否开启 SHOW VARIABLES LIKE 'log_bin'; -- 设置 binlog 格式为 ROW(持久化到配置文件) SET GLOBAL binlog_format = 'ROW';

- **Flink CDC 连接器 Jar 包**:下载对应版本的 Flink CDC 连接器(如 `flink-sql-connector-mysql-cdc-3.0.0.jar`),放入 Flink 的 `lib/` 目录下。

## 实战一:使用 Flink SQL 捕获 MySQL 变更
### 步骤 1:创建 MySQL 测试表并插入数据
在你的 MySQL 数据库中执行:
```sql
CREATE DATABASE inventory;
USE inventory;
CREATE TABLE products (
id INT PRIMARY KEY,
name VARCHAR(100),
price DECIMAL(10,2),
quantity INT
);
INSERT INTO products VALUES (1, 'Laptop', 1200.00, 10);
INSERT INTO products VALUES (2, 'Mouse', 25.00, 50);
# 进入 Flink 目录,启动 SQL 客户端
./bin/sql-client.sh

步骤 3:定义 MySQL CDC 源表

在 Flink SQL 客户端中执行:

CREATE TABLE products_source (
  id INT,
  name STRING,
  price DECIMAL(10, 2),
  quantity INT,
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'mysql-cdc',
  'hostname' = 'localhost',
  'port' = '3306',
  'username' = 'root',
  'password' = '123456',
  'database-name' = 'inventory',
  'table-name' = 'products',
  -- 可选:指定 server-id 避免冲突
  'server-id' = '5400-5404'
);

此时 Flink 即开始全量读取 products 表数据,并持续监听后续变更。

步骤 4:验证全量数据

SELECT * FROM products_source;

你会看到两行记录(id=1 和 id=2)。

步骤 5:观察增量数据

在 MySQL 中执行任意更新:

UPDATE products SET quantity = 15 WHERE id = 1;

回到 Flink SQL 客户端,再次查询 products_source 或通过持续查询观察变化:

SELECT * FROM products_source;

输出会显示更新后的 quantity 为 15。如果是 DELETE 操作,默认会输出一条空消息(取决于下游表定义)。

步骤 6:结果写出(Sink)

通常我们会将捕获的变更写入 Kafka、另一数据库或打印。这里以打印到控制台为例(仅用于调试):

CREATE TABLE console_sink (
  id INT,
  name STRING,
  price DECIMAL(10, 2),
  quantity INT,
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'print'
);

INSERT INTO console_sink SELECT * FROM products_source;

现在,任何对 MySQL products 表的变更都会实时打印到 Flink TaskManager 的日志中。

实战二:DataStream API 使用 MySQL CDC

如果需要在 Java/Scala 代码中更灵活地处理数据,可以使用 DataStream API。

Maven 依赖

<dependency>
  <groupId>com.ververica</groupId>
  <artifactId>flink-connector-mysql-cdc</artifactId>
  <version>3.0.0</version>
</dependency>

示例代码

import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.connectors.mysql.table.StartupOptions;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class MySqlCdcExample {
    public static void main(String[] args) throws Exception {
        MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
            .hostname("localhost")
            .port(3306)
            .databaseList("inventory")          // 支持正则
            .tableList("inventory.products")    // 支持正则
            .username("root")
            .password("123456")
            .startupOptions(StartupOptions.initial()) // 先全量,后增量
            .deserializer(new JsonDebeziumDeserializationSchema()) // 输出 JSON
            .build();

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL CDC Source")
           .print();

        env.execute("MySQL CDC Job");
    }
}

运行该程序,控制台会输出 JSON 格式的变更事件:

{"before":null,"after":{"id":1,"name":"Laptop","price":1200.0,"quantity":10},"source":{...},"op":"r"}
{"before":{"id":1,...},"after":{"id":1,"quantity":15},"op":"u"}

op 字段标识操作类型:r 表示 snapshot 读取,c 表示 insert,u 表示 update,d 表示 delete。

高级配置与最佳实践

启动模式

通过 startupOptions 控制首次启动行为:

  • initial(默认):全量读取后切换增量。
  • latest-offset:只读取最新的增量数据,跳过历史快照。
  • timestamp:从指定时间戳开始的 binlog 位置启动。

并行度与拆分全量快照

对于大表,全量读取可能成为瓶颈。可以在 MySQL 选项中设置:

'scan.incremental.snapshot.enabled' = 'true'
'scan.incremental.snapshot.chunk.size' = '8096'

开启增量快照后,Flink 会将表分成多个 chunk 并行读取,提升初始化速度。

处理 Schema 变更

Flink CDC 默认不会自动适应表结构变化(如加减列)。对于简单的兼容变更(加列且保证下游能为新字段提供默认值),可以设置:

'include-schema-change' = 'true'

但需要下游支持 schema 演化。更稳妥的做法是规划好 Schema 并通过 CDC 配合 Schema Registry 使用。

Exactly-Once 保障

Flink CDC 结合 checkpoint 机制可以提供端到端 exactly-once。确保启用 checkpoint:

env.enableCheckpointing(5000);

对于 Sink 连接器,需选择支持 exactly-once 的(如 Kafka、JDBC with XA)。

监控与运维

  • 观察 Source 的 metrics:全量阶段进度、Binlog 延迟等。
  • 设置合理的 server-id 范围,避免多任务冲突。
  • 数据库用户需要 RELOAD, SUPER, REPLICATION SLAVE, REPLICATION CLIENT 等权限。

常见问题排查

Q:启动报错 Access denied; you need (at least one of) the SUPER, REPLICATION CLIENT privilege(s)
A:授予用户对应权限:

GRANT SUPER, RELOAD, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'user'@'%';