Flink CDC 实时捕获数据库变更
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);
步骤 2:启动 Flink SQL 客户端
# 进入 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'@'%';