Flink CDC 实战案例详解
基础概念
Flink CDC(Change Data Capture)是Apache Flink提供的一个强大的数据捕获框架,能够实时捕获数据库中的数据变更(插入、更新、删除),并将其转换为流式事件进行处理。

经典案例:MySQL实时同步到Kafka
场景描述
将MySQL中订单表的变更实时同步到Kafka,供下游系统消费。
环境准备
<!-- pom.xml 依赖配置 -->
<dependencies>
<!-- Flink CDC 依赖 -->
<dependency>
<groupId>com.ververica</groupId>
<artifactId>flink-connector-mysql-cdc</artifactId>
<version>2.4.0</version>
</dependency>
<!-- Flink Kafka 连接器 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>1.17.0</version>
</dependency>
<!-- Flink Stream API -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.17.0</version>
</dependency>
</dependencies>
MySQL数据准备
-- 创建数据库
CREATE DATABASE flink_cdc_demo;
USE flink_cdc_demo;
-- 创建订单表
CREATE TABLE orders (
id INT PRIMARY KEY AUTO_INCREMENT,
order_no VARCHAR(50) NOT NULL,
user_id INT NOT NULL,
product_name VARCHAR(100),
amount DECIMAL(10,2),
status VARCHAR(20),
create_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
update_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
);
-- 插入测试数据
INSERT INTO orders (order_no, user_id, product_name, amount, status)
VALUES
('ORD-001', 1001, 'iPhone 15', 6999.00, 'CREATED'),
('ORD-002', 1002, 'MacBook Pro', 15999.00, 'PAID'),
('ORD-003', 1003, 'AirPods', 1299.00, 'SHIPPED');
-- 创建用户表
CREATE TABLE users (
id INT PRIMARY KEY AUTO_INCREMENT,
name VARCHAR(50),
email VARCHAR(100)
);
INSERT INTO users (name, email) VALUES ('张三', 'zhangsan@example.com');
核心代码实现
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.util.Collector;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.connector.kafka.sink.KafkaSink;
import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema;
import org.apache.flink.connector.kafka.sink.KafkaSinkBuilder;
public class MySqlCDCToKafka {
public static void main(String[] args) throws Exception {
// 1. 创建执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// 启用Checkpoint
env.enableCheckpointing(5000);
// 2. 创建MySQL CDC Source
MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("flink_cdc_demo") // 监控的数据库
.tableList("flink_cdc_demo.orders") // 监控的表
.username("root")
.password("password")
.deserializer(new JsonDebeziumDeserializationSchema()) // 转换为JSON
.startupOptions(StartupOptions.initial()) // 启动模式:初始化
.build();
// 3. 创建Kafka Sink
KafkaSink<String> kafkaSink = KafkaSink.<String>builder()
.setBootstrapServers("localhost:9092")
.setRecordSerializer(KafkaRecordSerializationSchema.builder()
.setTopic("cdc-orders")
.setValueSerializationSchema(new SimpleStringSchema())
.build()
)
.setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
.build();
// 4. 构建数据流
DataStream<String> stream = env
.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL CDC Source")
.name("MySQL CDC Source");
// 5. 数据处理(可选:解析和转换)
DataStream<String> processedStream = stream
.process(new ProcessFunction<String, String>() {
@Override
public void processElement(String value, Context ctx, Collector<String> out) throws Exception {
// 这里可以添加自定义处理逻辑
// 数据清洗、格式转换、业务逻辑处理等
out.collect(value);
}
})
.name("Data Processing");
// 6. 写入Kafka
processedStream.sinkTo(kafkaSink).name("Kafka Sink");
// 7. 执行任务
env.execute("MySQL CDC to Kafka");
}
}
配置其他启动模式
// 多种启动模式配置
// 1. 初始化模式(默认):先读取现有数据,然后继续读取变更
StartupOptions.initial()
// 2. 最早模式:从最早的binlog开始读取
StartupOptions.earliest()
// 3. 最新模式:只从当前时间开始读取变更
StartupOptions.latest()
// 4. 指定时间戳
StartupOptions.timestamp(1700000000000L)
// 5. 指定偏移量
StartupOptions.specificOffset("mysql-bin.000001", 4L, 12345L)
进阶案例:多表关联与状态管理
import org.apache.flink.api.common.state.MapState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction;
import org.apache.flink.util.Collector;
public class MultiTableJoinCDC {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// 1. 创建订单CDC Source
MySqlSource<String> orderSource = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("flink_cdc_demo")
.tableList("flink_cdc_demo.orders")
.username("root")
.password("password")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
// 2. 创建用户CDC Source
MySqlSource<String> userSource = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("flink_cdc_demo")
.tableList("flink_cdc_demo.users")
.username("root")
.password("password")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
// 3. 获取数据流
DataStream<String> orderStream = env.fromSource(
orderSource,
WatermarkStrategy.noWatermarks(),
"Order CDC Source"
);
DataStream<String> userStream = env.fromSource(
userSource,
WatermarkStrategy.noWatermarks(),
"User CDC Source"
);
// 4. 使用Connect合并两个流并关联
DataStream<String> result = orderStream
.connect(userStream)
.process(new OrderUserJoinFunction())
.name("Order-User Join");
// 5. 输出结果
result.print();
env.execute("Multi-table Join CDC");
}
// 自定义连接函数
public static class OrderUserJoinFunction
extends KeyedBroadcastProcessFunction<String, String, String, String> {
private MapState<String, String> userState;
@Override
public void open(Configuration parameters) {
MapStateDescriptor<String, String> userDescriptor =
new MapStateDescriptor<>("users", String.class, String.class);
userState = getRuntimeContext().getMapState(userDescriptor);
}
@Override
public void processElement(String orderJson, ReadOnlyContext ctx,
Collector<String> out) throws Exception {
// 解析订单JSON,获取user_id
String userId = extractUserId(orderJson);
String userInfo = userState.get(userId);
if (userInfo != null) {
// 关联用户信息
String enrichedOrder = enrichOrderWithUser(orderJson, userInfo);
out.collect(enrichedOrder);
}
}
@Override
public void processBroadcastElement(String userJson, Context ctx,
Collector<String> out) throws Exception {
// 更新用户状态
String userId = extractUserId(userJson);
userState.put(userId, userJson);
}
}
}
实际生产案例:实时数仓同步
// 完整的生产级配置示例
public class ProductionCDCApplication {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 生产环境配置
Configuration config = new Configuration();
config.setInteger("taskmanager.numberOfTaskSlots", 4);
env.configure(config);
// 启用Checkpoint
CheckpointConfig checkpointConfig = env.getCheckpointConfig();
checkpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
checkpointConfig.setCheckpointInterval(60000); // 1分钟
checkpointConfig.setCheckpointTimeout(60000);
checkpointConfig.setMaxConcurrentCheckpoints(1);
checkpointConfig.setMinPauseBetweenCheckpoints(5000);
// 设置状态后端
env.setStateBackend(new HashMapStateBackend());
env.getCheckpointConfig().setCheckpointStorage("hdfs://namenode:8020/flink/checkpoints");
// 创建多表监控的CDC Source
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("mysql-master")
.port(3306)
.databaseList("business_db") // 业务数据库
.tableList("business_db.orders,business_db.payments,business_db.shipments") // 多表
.username("cdc_user")
.password("cdc_password")
.serverId("5400-5404") // 分配多个server ID用于并行读取
.debeziumProperties(new HashMap<String, String>() {{
put("snapshot.locking.mode", "none"); // 无锁快照
put("database.server.name", "business_db");
put("include.schema.changes", "true");
}})
.deserializer(new JsonDebeziumDeserializationSchema())
.startupOptions(StartupOptions.latest())
.build();
DataStream<String> stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
"Production CDC Source"
)
.setParallelism(4) // 设置并行度
.rebalance();
// 写入多个目标
// 1. 写入Kafka
stream.sinkTo(createKafkaSink("cdc-topic"));
// 2. 写入StarRocks/ClickHouse
stream.addSink(createStarRocksSink());
// 3. 本地调试输出
if (isDebugMode()) {
stream.print();
}
env.execute("Production CDC Sync Job");
}
}
常见问题与解决方案
数据一致性
// 确保至少一次语义 CheckpointConfig config = env.getCheckpointConfig(); config.setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE); // 或精确一次性 config.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
性能优化
# flink-conf.yaml 配置 # 增加Socket接收缓冲区 taskmanager.network.memory.buffer-debloat.enabled: true # 优化网络缓冲区 taskmanager.memory.network.min: 64mb taskmanager.memory.network.max: 128mb # 并行度配置 parallelism.default: 4
监控告警
// 添加Metrics监控
stream
.map(new RichMapFunction<String, String>() {
private transient Counter counter;
@Override
public void open(Configuration parameters) {
counter = getRuntimeContext()
.getMetricGroup()
.counter("cdc_records_count");
}
@Override
public String map(String value) throws Exception {
counter.inc();
return value;
}
});
CDC数据格式示例
// MySQL CDC 输出的JSON格式
{
"before": {
"id": 1,
"order_no": "ORD-001",
"status": "CREATED"
},
"after": {
"id": 1,
"order_no": "ORD-001",
"status": "PAID"
},
"source": {
"db": "flink_cdc_demo",
"table": "orders",
"server_id": 123,
"ts_sec": 1700000000
},
"op": "u", // c:创建, u:更新, d:删除, r:快照读取
"ts_ms": 1700000000000
}
最佳实践总结
-
启动策略选择
- 首次使用:
initial模式全量+增量 - 已有offset:
latest模式只读增量
- 首次使用:
-
性能调优
- 增大
server.id分配范围 - 调整
debezium.snapshot.fetch.size - 启用并行读取
- 增大
-
数据质量
- 实现数据校验机制
- 添加重试和死信队列
- 监控延迟和吞吐量
-
运维建议
- 定期检查binlog保留时间
- 监控MySQL性能指标
- 做好Flink状态备份
这个案例涵盖了从基础到生产的Flink CDC应用,可以根据实际需求选择合适的实现方案。