概述
Flink CDC (Change Data Capture) 是基于 Apache Flink 的数据变更捕获框架,能够实时捕获数据库的增删改操作,并将变更事件流式传输到下游系统。它提供了一站式的实时数据同步解决方案。
核心特性
- 全量 + 增量同步: 支持历史数据全量同步和实时增量同步
- Exactly-Once 语义: 保证数据不丢失、不重复
- 无锁读取: 不影响在线业务,对源库压力小
- 分布式架构: 支持水平扩展,处理大规模数据
- Schema 演进: 自动感知并同步表结构变更
- 断点续传: 支持从中断位置继续同步
应用场景
| 场景 | 描述 | 优势 |
|---|---|---|
| 实时数据同步 | 数据库之间的实时同步 | 低延迟、高可靠 |
| 实时数仓构建 | 构建实时数据仓库 | 实时 ETL、数据清洗 |
| 缓存更新 | 数据库到缓存的同步 | 自动失效、实时更新 |
| 搜索索引构建 | 数据库到搜索引擎同步 | 近实时搜索 |
| 微服务数据同步 | 跨服务数据共享 | 解耦服务、数据一致性 |
💡 提示: Flink CDC 相比传统的批量同步工具,具有更低的延迟和更小的资源开销。
CDC 基础概念
什么是 CDC
Change Data Capture(变更数据捕获)是一种用于识别和捕获数据库中数据变更的技术:
传统批量同步 vs CDC 实时同步
批量同步:
┌─────────┐ 定时全量 ┌─────────┐
│ Source │ ────────> │ Target │
│ DB │ (高延迟) │ DB │
└─────────┘ └─────────┘
CDC 同步:
┌─────────┐ 实时增量 ┌─────────┐
│ Source │ ────────> │ Target │
│ DB │ (低延迟) │ DB │
└─────────┘ └─────────┘
↓
Binlog/WAL
CDC 实现方式对比
| 方式 | 原理 | 优点 | 缺点 |
|---|---|---|---|
| 基于查询 | 定期查询变更 | 实现简单 | 延迟高、压力大 |
| 基于触发器 | 数据库触发器 | 实时性好 | 侵入性强、性能影响 |
| 基于日志 | 解析数据库日志 | 低侵入、低延迟 | 实现复杂 |
Flink CDC 采用基于日志的方式,通过读取数据库的事务日志(如 MySQL 的 binlog、PostgreSQL 的 WAL)来捕获数据变更。
Flink CDC 架构
整体架构
┌─────────────────────────────────────────────────────────────┐
│ Flink CDC 架构 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 数据源层 │
│ ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐ │
│ │ MySQL │ │ PG │ │ Oracle │ │MongoDB │ │
│ └───┬────┘ └───┬────┘ └───┬────┘ └───┬────┘ │
│ │ │ │ │ │
│ ▼ ▼ ▼ ▼ │
│ ┌─────────────────────────────────────────┐ │
│ │ CDC Source Connector │ │
│ │ ┌─────────┐ ┌─────────┐ ┌─────────┐ │ │
│ │ │Debezium │ │ Reader │ │Snapshot │ │ │
│ │ │ Engine │ │ Thread │ │ Thread │ │ │
│ │ └─────────┘ └─────────┘ └─────────┘ │ │
│ └─────────────────────────────────────────┘ │
│ │ │
│ ▼ │
│ ┌─────────────────────────────────────────┐ │
│ │ Flink Runtime │ │
│ │ ┌─────────┐ ┌─────────┐ ┌─────────┐ │ │
│ │ │ State │ │Checkpoint│ │Watermark│ │ │
│ │ │ Backend │ │ │ │ │ │ │
│ │ └─────────┘ └─────────┘ └─────────┘ │ │
│ └─────────────────────────────────────────┘ │
│ │ │
│ ▼ │
│ 数据处理层 │
│ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
│ │Transform│ │ Filter │ │ Route │ │
│ └─────────┘ └─────────┘ └─────────┘ │
│ │ │
│ ▼ │
│ 目标系统层 │
│ ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐ │
│ │ Kafka │ │ ES │ │ HBase │ │Iceberg │ │
│ └────────┘ └────────┘ └────────┘ └────────┘ │
│ │
└─────────────────────────────────────────────────────────────┘
核心组件
1. Debezium Engine
负责从数据库日志中读取变更事件:
- 管理数据库连接
- 解析二进制日志
- 转换为统一的事件格式
2. Source Connector
Flink 数据源连接器:
- 管理读取位置(offset)
- 协调全量和增量读取
- 处理故障恢复
3. Snapshot Reader
全量数据读取器:
- 并行读取历史数据
- 生成一致性快照
- 无锁读取实现
4. Stream Reader
增量数据读取器:
- 实时读取变更日志
- 维护读取位置
- 处理 Schema 变更
快速开始
环境准备
<!-- Maven 依赖 -->
<dependencies>
<!-- Flink CDC -->
<dependency>
<groupId>com.ververica</groupId>
<artifactId>flink-sql-connector-mysql-cdc</artifactId>
<version>2.4.0</version>
</dependency>
<!-- Flink 依赖 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.17.0</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java-bridge</artifactId>
<version>1.17.0</version>
</dependency>
</dependencies>
基础示例
1. DataStream API 方式
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 {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 设置并行度
env.setParallelism(4);
// 创建 MySQL CDC Source
MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("test_db") // 监控的数据库
.tableList("test_db.users", "test_db.orders") // 监控的表
.username("root")
.password("password")
.startupOptions(StartupOptions.initial()) // 从最早位置开始读取
.deserializer(new JsonDebeziumDeserializationSchema()) // JSON 格式
.build();
// 使用 CDC Source
env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
.print();
env.execute("MySQL CDC Example");
}
}
2. Table API 方式
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
public class MySqlCDCTableExample {
public static void main(String[] args) {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// 创建 CDC 源表
tableEnv.executeSql(
"CREATE TABLE users_source (" +
" id INT," +
" name STRING," +
" age INT," +
" email STRING," +
" update_time TIMESTAMP(3)," +
" PRIMARY KEY (id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'mysql-cdc'," +
" 'hostname' = 'localhost'," +
" 'port' = '3306'," +
" 'username' = 'root'," +
" 'password' = 'password'," +
" 'database-name' = 'test_db'," +
" 'table-name' = 'users'" +
")"
);
// 创建目标表
tableEnv.executeSql(
"CREATE TABLE users_sink (" +
" id INT," +
" name STRING," +
" age INT," +
" email STRING," +
" update_time TIMESTAMP(3)," +
" PRIMARY KEY (id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'elasticsearch-7'," +
" 'hosts' = 'http://localhost:9200'," +
" 'index' = 'users'" +
")"
);
// 同步数据
tableEnv.executeSql("INSERT INTO users_sink SELECT * FROM users_source");
}
}
支持的数据源
MySQL CDC
1. 前置条件
-- 检查 binlog 是否开启
SHOW VARIABLES LIKE 'log_bin';
-- 查看 binlog 格式(需要是 ROW)
SHOW VARIABLES LIKE 'binlog_format';
-- 配置 my.cnf
[mysqld]
log-bin=mysql-bin
binlog-format=ROW
binlog-row-image=FULL
expire_logs_days=7
2. 高级配置
public class MySqlCDCAdvanced {
public static MySqlSource<String> createAdvancedMySqlSource() {
Properties jdbcProperties = new Properties();
jdbcProperties.put("useSSL", "false");
jdbcProperties.put("allowPublicKeyRetrieval", "true");
Properties debeziumProperties = new Properties();
// 设置时区
debeziumProperties.put("database.serverTimezone", "Asia/Shanghai");
// 处理 decimal 类型
debeziumProperties.put("decimal.handling.mode", "string");
return MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("db1", "db2") // 多数据库
.tableList("db1.table1", "db2.*") // 支持正则
.username("root")
.password("password")
.serverId("5400-5404") // 服务器 ID 范围
.splitSize(8096) // 全量读取时的分片大小
.fetchSize(1024) // JDBC fetch size
.connectTimeout(Duration.ofSeconds(30))
.connectionPoolSize(20) // 连接池大小
.jdbcProperties(jdbcProperties)
.debeziumProperties(debeziumProperties)
.startupOptions(StartupOptions.specificOffset("mysql-bin.000003", 4))
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
}
}
3. 自定义反序列化
public class CustomDeserializationSchema implements DebeziumDeserializationSchema<RowData> {
@Override
public void deserialize(SourceRecord record, Collector<RowData> out) throws Exception {
// 获取操作类型
Envelope.Operation op = Envelope.operationFor(record);
// 获取数据
Struct value = (Struct) record.value();
Struct source = value.getStruct("source");
String database = source.getString("db");
String table = source.getString("table");
// 根据操作类型处理
switch (op) {
case CREATE:
case READ:
Struct after = value.getStruct("after");
RowData insertRow = extractRowData(after, RowKind.INSERT);
out.collect(insertRow);
break;
case UPDATE:
Struct before = value.getStruct("before");
Struct afterUpdate = value.getStruct("after");
// 发送删除和插入事件(CDC 语义)
if (before != null) {
RowData deleteRow = extractRowData(before, RowKind.DELETE);
out.collect(deleteRow);
}
RowData insertRowUpdate = extractRowData(afterUpdate, RowKind.INSERT);
out.collect(insertRowUpdate);
break;
case DELETE:
Struct beforeDelete = value.getStruct("before");
RowData deleteRow = extractRowData(beforeDelete, RowKind.DELETE);
out.collect(deleteRow);
break;
}
}
private RowData extractRowData(Struct struct, RowKind rowKind) {
// 实现数据提取逻辑
// ...
}
@Override
public TypeInformation<RowData> getProducedType() {
return TypeInformation.of(RowData.class);
}
}
PostgreSQL CDC
1. 配置要求
-- postgresql.conf
wal_level = logical
max_replication_slots = 4
max_wal_senders = 4
-- 创建发布
CREATE PUBLICATION flink_publication FOR ALL TABLES;
-- 创建复制槽
SELECT * FROM pg_create_logical_replication_slot('flink_slot', 'pgoutput');
2. 使用示例
public class PostgreSQLCDCExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
PostgreSQLSource<String> postgresSource = PostgreSQLSource.<String>builder()
.hostname("localhost")
.port(5432)
.database("postgres")
.schemaList("public")
.tableList("public.users", "public.orders")
.username("postgres")
.password("password")
.slotName("flink_slot")
.decodingPluginName("pgoutput") // 或 wal2json
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
env.fromSource(postgresSource, WatermarkStrategy.noWatermarks(), "PostgreSQL Source")
.print();
env.execute("PostgreSQL CDC Example");
}
}
Oracle CDC
1. 配置 Oracle
-- 启用归档日志
ALTER DATABASE ARCHIVELOG;
-- 启用补充日志
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;
-- 创建用户并授权
CREATE USER flinkuser IDENTIFIED BY flinkpw;
GRANT CREATE SESSION TO flinkuser;
GRANT SELECT ANY TABLE TO flinkuser;
GRANT SELECT_CATALOG_ROLE TO flinkuser;
GRANT EXECUTE_CATALOG_ROLE TO flinkuser;
GRANT SELECT ANY TRANSACTION TO flinkuser;
GRANT LOGMINING TO flinkuser;
2. 使用示例
public class OracleCDCExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
OracleSource<String> oracleSource = OracleSource.<String>builder()
.hostname("localhost")
.port(1521)
.database("ORCL") // SID
.schemaList("FLINKUSER")
.tableList("FLINKUSER.USERS")
.username("flinkuser")
.password("flinkpw")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
env.fromSource(oracleSource, WatermarkStrategy.noWatermarks(), "Oracle Source")
.print();
env.execute("Oracle CDC Example");
}
}
MongoDB CDC
1. 前置条件
// MongoDB 必须是副本集模式
// 检查副本集状态
rs.status()
// 如果是单机,转换为副本集
rs.initiate()
2. 使用示例
public class MongoDBCDCExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
MongoDBSource<String> mongoSource = MongoDBSource.<String>builder()
.hosts("localhost:27017")
.username("flinkuser")
.password("flinkpw")
.databaseList("test") // 支持正则
.collectionList("test.users", "test.orders")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
env.fromSource(mongoSource, WatermarkStrategy.noWatermarks(), "MongoDB Source")
.print();
env.execute("MongoDB CDC Example");
}
}
SQL Server CDC
1. 启用 CDC
-- 启用数据库级别 CDC
EXEC sys.sp_cdc_enable_db;
-- 启用表级别 CDC
EXEC sys.sp_cdc_enable_table
@source_schema = N'dbo',
@source_name = N'users',
@role_name = NULL;
2. 使用示例
public class SqlServerCDCExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
SqlServerSource<String> sqlServerSource = SqlServerSource.<String>builder()
.hostname("localhost")
.port(1433)
.database("testdb")
.tableList("dbo.users")
.username("sa")
.password("Password!")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
env.fromSource(sqlServerSource, WatermarkStrategy.noWatermarks(), "SQL Server Source")
.print();
env.execute("SQL Server CDC Example");
}
}
高级特性
全量与增量同步
1. 启动模式配置
public class StartupOptionsExample {
public static void demonstrateStartupOptions() {
// 1. 从最早位置开始(包括全量)
StartupOptions initial = StartupOptions.initial();
// 2. 从最新位置开始(仅增量)
StartupOptions latest = StartupOptions.latest();
// 3. 从指定位置开始
StartupOptions specificOffset = StartupOptions.specificOffset("mysql-bin.000003", 4);
// 4. 从指定时间戳开始
StartupOptions timestamp = StartupOptions.timestamp(1000L);
// 5. 跳过全量阶段,从变更流开始
StartupOptions snapshot = StartupOptions.snapshot();
}
}
2. 分片并行读取
public class ParallelSnapshotReading {
public static MySqlSource<String> createParallelSource() {
return MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("large_db")
.tableList("large_db.huge_table")
.username("root")
.password("password")
.splitSize(8096) // 每个分片的大小
.splitMetaGroupSize(1000) // 分片元数据组大小
.fetchSize(1024) // JDBC 批量大小
.connectTimeout(Duration.ofSeconds(30))
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
}
}
Schema 演进
1. 自动 Schema 变更感知
public class SchemaEvolutionExample {
public static void handleSchemaEvolution() {
// 配置 Debezium 属性以处理 Schema 变更
Properties debeziumProperties = new Properties();
// 包含 Schema 变更事件
debeziumProperties.put("include.schema.changes", "true");
// 存储 Schema 历史
debeziumProperties.put("database.history.store.only.monitored.tables.ddl", "true");
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("test")
.tableList("test.*")
.username("root")
.password("password")
.debeziumProperties(debeziumProperties)
.deserializer(new SchemaAwareDeserializer())
.build();
}
}
// 自定义 Schema 感知反序列化器
public class SchemaAwareDeserializer implements DebeziumDeserializationSchema<String> {
@Override
public void deserialize(SourceRecord record, Collector<String> out) throws Exception {
// 检查是否是 Schema 变更事件
if (record.topic().endsWith(".DDL")) {
handleSchemaChange(record);
return;
}
// 处理数据变更事件
String result = convertToJson(record);
out.collect(result);
}
private void handleSchemaChange(SourceRecord record) {
// 处理 DDL 变更
Struct value = (Struct) record.value();
String ddl = value.getString("ddl");
System.out.println("Schema change detected: " + ddl);
// 可以将 Schema 变更信息发送到元数据管理系统
}
}
分布式快照
1. 一致性保证
public class ConsistentSnapshotExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 启用 Checkpoint 保证一致性
env.enableCheckpointing(10000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("test")
.tableList("test.*")
.username("root")
.password("password")
// 使用一致性快照读取
.startupOptions(StartupOptions.initial())
.includeSchemaChanges(true) // 包含 schema 变更
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
DataStream<String> stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
"MySQL Source"
);
// 处理数据
stream.process(new ProcessFunction<String, String>() {
private transient ValueState<Long> offsetState;
@Override
public void open(Configuration parameters) throws Exception {
ValueStateDescriptor<Long> descriptor =
new ValueStateDescriptor<>("offset", Long.class);
offsetState = getRuntimeContext().getState(descriptor);
}
@Override
public void processElement(String value, Context ctx, Collector<String> out)
throws Exception {
// 保存 offset 到状态
JSONObject json = JSON.parseObject(value);
Long offset = json.getLong("pos");
if (offset != null) {
offsetState.update(offset);
}
out.collect(value);
}
});
env.execute("Consistent Snapshot Example");
}
}
Exactly-Once 语义
1. 端到端 Exactly-Once
public class ExactlyOnceExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 配置 Exactly-Once
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.enableCheckpointing(5000);
// MySQL CDC Source
MySqlSource<RowData> mySqlSource = createMySqlSource();
DataStream<RowData> mysqlStream = env.fromSource(
mySqlSource,
WatermarkStrategy.noWatermarks(),
"MySQL Source"
);
// Kafka Sink with Exactly-Once
KafkaSink<RowData> kafkaSink = KafkaSink.<RowData>builder()
.setBootstrapServers("localhost:9092")
.setRecordSerializer(new RowDataSerializationSchema())
.setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix("flink-cdc-")
.build();
mysqlStream.sinkTo(kafkaSink);
env.execute("Exactly-Once CDC Pipeline");
}
}
性能优化
并行度优化
1. Source 并行度设置
public class ParallelismOptimization {
public static void optimizeParallelism(StreamExecutionEnvironment env) {
// 全局并行度
env.setParallelism(4);
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("test")
.tableList("test.large_table")
.username("root")
.password("password")
// 服务器 ID 范围,支持并行读取
.serverId("5400-5403") // 4个并行度
.splitSize(8096) // 合理的分片大小
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
// Source 并行度
DataStream<String> stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
"MySQL Source"
).setParallelism(4);
// 下游处理并行度可以不同
stream.map(new MyMapper()).setParallelism(8)
.keyBy(new MyKeySelector()).setParallelism(16)
.sink(new MySink()).setParallelism(4);
}
}
2. 表分片策略
public class TableSplitStrategy {
public static class CustomSplitStrategy {
// 根据主键范围分片
public List<MySqlSplit> splitTable(TableId tableId, Object[] primaryKeys) {
List<MySqlSplit> splits = new ArrayList<>();
// 获取主键范围
Object minKey = getMinPrimaryKey(tableId);
Object maxKey = getMaxPrimaryKey(tableId);
// 计算分片
long range = (Long)maxKey - (Long)minKey;
long splitSize = range / 4; // 4个分片
for (int i = 0; i < 4; i++) {
long start = (Long)minKey + i * splitSize;
long end = (i == 3) ? (Long)maxKey : start + splitSize;
MySqlSplit split = new MySqlSplit(
tableId,
"split-" + i,
start,
end
);
splits.add(split);
}
return splits;
}
}
}
内存优化
1. 缓冲区配置
public class MemoryOptimization {
public static void configureMemory() {
Properties debeziumProperties = new Properties();
// Kafka 缓冲区设置
debeziumProperties.put("max.queue.size", "16384");
debeziumProperties.put("max.batch.size", "8192");
// 数据库连接池
Properties jdbcProperties = new Properties();
jdbcProperties.put("connectionPoolSize", "5");
jdbcProperties.put("connectTimeout", "30");
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("test")
.tableList("test.*")
.username("root")
.password("password")
.fetchSize(1000) // JDBC fetch size
.connectionPoolSize(5) // 连接池大小
.jdbcProperties(jdbcProperties)
.debeziumProperties(debeziumProperties)
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
}
}
2. 状态管理优化
public class StateOptimization {
public static void optimizeState(StreamExecutionEnvironment env) {
// 使用 RocksDB 状态后端
env.setStateBackend(new RocksDBStateBackend("hdfs://checkpoints"));
// 增量 Checkpoint
env.getCheckpointConfig().enableIncrementalCheckpoints();
// 状态 TTL
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.hours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
// 应用到状态描述符
ValueStateDescriptor<String> descriptor =
new ValueStateDescriptor<>("offset-state", String.class);
descriptor.enableTimeToLive(ttlConfig);
}
}
网络优化
1. 批量处理
public class BatchProcessingOptimization {
public static class BatchingProcessFunction
extends ProcessFunction<String, List<String>> {
private transient ListState<String> bufferedElements;
private final int batchSize = 1000;
private final long batchInterval = 5000L; // 5秒
@Override
public void open(Configuration parameters) {
ListStateDescriptor<String> descriptor =
new ListStateDescriptor<>("buffer", String.class);
bufferedElements = getRuntimeContext().getListState(descriptor);
}
@Override
public void processElement(String value, Context ctx,
Collector<List<String>> out) throws Exception {
bufferedElements.add(value);
// 注册定时器
ctx.timerService().registerProcessingTimeTimer(
ctx.timerService().currentProcessingTime() + batchInterval
);
// 检查批次大小
List<String> batch = new ArrayList<>();
for (String element : bufferedElements.get()) {
batch.add(element);
}
if (batch.size() >= batchSize) {
out.collect(batch);
bufferedElements.clear();
}
}
@Override
public void onTimer(long timestamp, OnTimerContext ctx,
Collector<List<String>> out) throws Exception {
List<String> batch = new ArrayList<>();
for (String element : bufferedElements.get()) {
batch.add(element);
}
if (!batch.isEmpty()) {
out.collect(batch);
bufferedElements.clear();
}
}
}
}
实战案例
MySQL 到 Elasticsearch
1. 完整示例
public class MySQL2ElasticsearchExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// 设置检查点
env.enableCheckpointing(5000);
// 创建 MySQL CDC 源表
tableEnv.executeSql(
"CREATE TABLE mysql_users (" +
" id BIGINT," +
" name STRING," +
" age INT," +
" email STRING," +
" address STRING," +
" update_time TIMESTAMP(3)," +
" PRIMARY KEY (id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'mysql-cdc'," +
" 'hostname' = 'localhost'," +
" 'port' = '3306'," +
" 'username' = 'root'," +
" 'password' = 'password'," +
" 'database-name' = 'test'," +
" 'table-name' = 'users'," +
" 'server-id' = '5400-5404'," +
" 'scan.startup.mode' = 'initial'" +
")"
);
// 创建 Elasticsearch 目标表
tableEnv.executeSql(
"CREATE TABLE es_users (" +
" id BIGINT," +
" name STRING," +
" age INT," +
" email STRING," +
" address STRING," +
" update_time TIMESTAMP(3)," +
" PRIMARY KEY (id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'elasticsearch-7'," +
" 'hosts' = 'http://localhost:9200'," +
" 'index' = 'users'," +
" 'document-id.key-delimiter' = '_'," +
" 'sink.bulk-flush.max-size' = '1000'," +
" 'sink.bulk-flush.max-actions' = '1000'," +
" 'sink.bulk-flush.interval' = '1s'," +
" 'format' = 'json'" +
")"
);
// 数据同步与转换
tableEnv.executeSql(
"INSERT INTO es_users " +
"SELECT " +
" id," +
" UPPER(name) as name," + // 转换为大写
" age," +
" email," +
" address," +
" update_time " +
"FROM mysql_users " +
"WHERE age >= 18" // 只同步成年用户
);
}
}
2. 自定义处理逻辑
public class CustomProcessingExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("test")
.tableList("test.orders")
.username("root")
.password("password")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
DataStream<String> mysqlStream = env.fromSource(
mySqlSource,
WatermarkStrategy.noWatermarks(),
"MySQL Source"
);
// 解析并处理数据
DataStream<Order> orderStream = mysqlStream
.map(json -> JSON.parseObject(json, Order.class))
.keyBy(Order::getUserId)
.window(TumblingProcessingTimeWindows.of(Time.minutes(5)))
.aggregate(new OrderAggregator());
// 写入 Elasticsearch
orderStream.addSink(
new ElasticsearchSink.Builder<>(
httpHosts,
new ElasticsearchSinkFunction<Order>() {
@Override
public void process(Order order, RuntimeContext ctx,
RequestIndexer indexer) {
IndexRequest request = Requests.indexRequest()
.index("order_stats")
.id(order.getUserId())
.source(JSON.toJSONString(order));
indexer.add(request);
}
}
).build()
);
env.execute("Custom MySQL to Elasticsearch Pipeline");
}
}
异构数据库同步
1. MySQL 到 PostgreSQL
public class MySQL2PostgreSQLSync {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// MySQL 源
tableEnv.executeSql(
"CREATE TABLE mysql_source (" +
" id INT," +
" name STRING," +
" create_time TIMESTAMP(3)," +
" PRIMARY KEY (id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'mysql-cdc'," +
" 'hostname' = 'mysql-host'," +
" 'port' = '3306'," +
" 'username' = 'root'," +
" 'password' = 'password'," +
" 'database-name' = 'source_db'," +
" 'table-name' = 'source_table'" +
")"
);
// PostgreSQL 目标
tableEnv.executeSql(
"CREATE TABLE pg_sink (" +
" id INT," +
" name STRING," +
" create_time TIMESTAMP(3)," +
" sync_time TIMESTAMP(3)," + // 添加同步时间
" PRIMARY KEY (id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'jdbc'," +
" 'url' = 'jdbc:postgresql://pg-host:5432/target_db'," +
" 'table-name' = 'target_table'," +
" 'username' = 'postgres'," +
" 'password' = 'password'," +
" 'sink.buffer-flush.max-rows' = '1000'," +
" 'sink.buffer-flush.interval' = '1s'" +
")"
);
// 同步数据
tableEnv.executeSql(
"INSERT INTO pg_sink " +
"SELECT " +
" id," +
" name," +
" create_time," +
" CURRENT_TIMESTAMP as sync_time " +
"FROM mysql_source"
);
}
}
2. 多源合并同步
public class MultiSourceMergeExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// MySQL 源
DataStream<UserEvent> mysqlStream = env.fromSource(
createMySqlSource(),
WatermarkStrategy.noWatermarks(),
"MySQL Source"
).map(new MySqlEventParser());
// PostgreSQL 源
DataStream<UserEvent> pgStream = env.fromSource(
createPostgreSQLSource(),
WatermarkStrategy.noWatermarks(),
"PostgreSQL Source"
).map(new PostgreSQLEventParser());
// MongoDB 源
DataStream<UserEvent> mongoStream = env.fromSource(
createMongoDBSource(),
WatermarkStrategy.noWatermarks(),
"MongoDB Source"
).map(new MongoEventParser());
// 合并多个数据源
DataStream<UserEvent> mergedStream = mysqlStream
.union(pgStream)
.union(mongoStream)
.keyBy(UserEvent::getUserId)
.process(new DeduplicationFunction()); // 去重
// 统一写入目标系统
mergedStream.addSink(createUnifiedSink());
env.execute("Multi-Source Merge Sync");
}
}
实时数据湖构建
1. CDC 到 Iceberg
public class CDC2IcebergExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// 配置 Iceberg Catalog
tableEnv.executeSql(
"CREATE CATALOG iceberg_catalog WITH (" +
" 'type' = 'iceberg'," +
" 'catalog-type' = 'hive'," +
" 'uri' = 'thrift://localhost:9083'," +
" 'warehouse' = 'hdfs://localhost:9000/user/hive/warehouse'" +
")"
);
tableEnv.useCatalog("iceberg_catalog");
// 创建 MySQL CDC 源
tableEnv.executeSql(
"CREATE TABLE mysql_orders (" +
" order_id BIGINT," +
" user_id BIGINT," +
" product_id BIGINT," +
" amount DECIMAL(10, 2)," +
" order_time TIMESTAMP(3)," +
" status STRING," +
" PRIMARY KEY (order_id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'mysql-cdc'," +
" 'hostname' = 'localhost'," +
" 'port' = '3306'," +
" 'username' = 'root'," +
" 'password' = 'password'," +
" 'database-name' = 'shop'," +
" 'table-name' = 'orders'" +
")"
);
// 创建 Iceberg 表
tableEnv.executeSql(
"CREATE TABLE IF NOT EXISTS iceberg_orders (" +
" order_id BIGINT," +
" user_id BIGINT," +
" product_id BIGINT," +
" amount DECIMAL(10, 2)," +
" order_time TIMESTAMP(3)," +
" status STRING," +
" dt STRING," + // 分区字段
" PRIMARY KEY (order_id) NOT ENFORCED" +
") PARTITIONED BY (dt) WITH (" +
" 'format-version' = '2'," +
" 'write.metadata.delete-after-commit.enabled' = 'true'," +
" 'write.metadata.previous-versions-max' = '3'" +
")"
);
// 实时同步到 Iceberg
tableEnv.executeSql(
"INSERT INTO iceberg_orders " +
"SELECT " +
" order_id," +
" user_id," +
" product_id," +
" amount," +
" order_time," +
" status," +
" DATE_FORMAT(order_time, 'yyyy-MM-dd') as dt " +
"FROM mysql_orders"
);
}
}
2. 实时数仓分层
public class RealtimeDataWarehouse {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// ODS 层 - 原始数据
createODSLayer(tableEnv);
// DWD 层 - 明细数据
createDWDLayer(tableEnv);
// DWS 层 - 汇总数据
createDWSLayer(tableEnv);
// ADS 层 - 应用数据
createADSLayer(tableEnv);
env.execute("Realtime Data Warehouse");
}
private static void createODSLayer(StreamTableEnvironment tableEnv) {
// 订单表 CDC
tableEnv.executeSql(
"CREATE TABLE ods_orders (" +
" order_id BIGINT," +
" user_id BIGINT," +
" amount DECIMAL(10, 2)," +
" order_time TIMESTAMP(3)," +
" PRIMARY KEY (order_id) NOT ENFORCED" +
") WITH (" +
" 'connector' = 'mysql-cdc'," +
" 'hostname' = 'localhost'," +
" 'port' = '3306'," +
" 'username' = 'root'," +
" 'password' = 'password'," +
" 'database-name' = 'shop'," +
" 'table-name' = 'orders'" +
")"
);
}
private static void createDWDLayer(StreamTableEnvironment tableEnv) {
// 订单明细宽表
tableEnv.executeSql(
"CREATE VIEW dwd_order_detail AS " +
"SELECT " +
" o.order_id," +
" o.user_id," +
" u.user_name," +
" u.user_level," +
" p.product_name," +
" p.category," +
" o.amount," +
" o.order_time " +
"FROM ods_orders o " +
"JOIN ods_users u ON o.user_id = u.user_id " +
"JOIN ods_products p ON o.product_id = p.product_id"
);
}
}
生产部署
集群部署配置
1. Flink 集群配置
# flink-conf.yaml
jobmanager:
memory:
process.size: 2g
taskmanager:
memory:
process.size: 4g
numberOfTaskSlots: 4
# CDC 相关配置
execution.checkpointing.interval: 10s
execution.checkpointing.mode: EXACTLY_ONCE
state.backend: rocksdb
state.checkpoints.dir: hdfs://namenode:9000/flink-checkpoints
2. Docker Compose 部署
version: '3.8'
services:
jobmanager:
image: flink:1.17-scala_2.12
ports:
- "8081:8081"
command: jobmanager
environment:
- JOB_MANAGER_RPC_ADDRESS=jobmanager
volumes:
- ./conf:/opt/flink/conf
- ./jars:/opt/flink/lib
taskmanager:
image: flink:1.17-scala_2.12
depends_on:
- jobmanager
command: taskmanager
scale: 3
environment:
- JOB_MANAGER_RPC_ADDRESS=jobmanager
volumes:
- ./conf:/opt/flink/conf
- ./jars:/opt/flink/lib
mysql:
image: mysql:8.0
environment:
MYSQL_ROOT_PASSWORD: password
MYSQL_DATABASE: test
ports:
- "3306:3306"
command: --log-bin=mysql-bin --binlog-format=ROW
3. Kubernetes 部署
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
name: flink-cdc-cluster
spec:
image: flink:1.17-scala_2.12
flinkVersion: v1_17
flinkConfiguration:
taskmanager.numberOfTaskSlots: "4"
state.backend: rocksdb
state.checkpoints.dir: s3://flink-checkpoints
execution.checkpointing.interval: 10s
serviceAccount: flink
jobManager:
resource:
memory: "2048m"
cpu: 1
taskManager:
replicas: 3
resource:
memory: "4096m"
cpu: 2
job:
jarURI: local:///opt/flink/lib/flink-cdc-job.jar
parallelism: 4
upgradeMode: savepoint
savepointTriggerNonce: 0
监控和运维
监控指标
1. 自定义监控指标
public class CDCMetricsExample {
public static class MetricsProcessFunction
extends ProcessFunction<String, String> {
// 延迟监控
private transient Histogram eventLatency;
// 吞吐量监控
private transient Meter throughput;
// 错误计数
private transient Counter errorCounter;
// 表级别监控
private transient Map<String, Counter> tableCounters;
@Override
public void open(Configuration parameters) {
MetricGroup metricGroup = getRuntimeContext().getMetricGroup();
eventLatency = metricGroup.histogram("event_latency",
new DescriptiveStatisticsHistogram(1000));
throughput = metricGroup.meter("throughput", new MeterView(60));
errorCounter = metricGroup.counter("errors");
tableCounters = new HashMap<>();
}
@Override
public void processElement(String value, Context ctx,
Collector<String> out) throws Exception {
try {
JSONObject event = JSON.parseObject(value);
// 计算延迟
long eventTime = event.getLong("ts_ms");
long currentTime = System.currentTimeMillis();
long latency = currentTime - eventTime;
eventLatency.update(latency);
// 记录吞吐量
throughput.markEvent();
// 表级别统计
String table = event.getString("table");
tableCounters.computeIfAbsent(table, k ->
getRuntimeContext().getMetricGroup()
.counter("table_" + k + "_count")
).inc();
out.collect(value);
} catch (Exception e) {
errorCounter.inc();
LOG.error("Process error", e);
}
}
}
}
2. Prometheus 集成
# flink-conf.yaml
metrics.reporters: prom
metrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter
metrics.reporter.prom.port: 9249
# Prometheus 配置
scrape_configs:
- job_name: 'flink'
static_configs:
- targets: ['jobmanager:9249', 'taskmanager1:9249', 'taskmanager2:9249']
告警配置
public class AlertingExample {
public static class LagMonitor extends ProcessFunction<String, Alert> {
private final long lagThreshold = 60000; // 1分钟
@Override
public void processElement(String value, Context ctx,
Collector<Alert> out) throws Exception {
JSONObject event = JSON.parseObject(value);
long eventTime = event.getLong("ts_ms");
long currentTime = System.currentTimeMillis();
long lag = currentTime - eventTime;
if (lag > lagThreshold) {
Alert alert = new Alert(
"HIGH_LAG",
"CDC lag exceeds threshold: " + lag + "ms",
Alert.Severity.WARNING
);
out.collect(alert);
// 发送告警
sendAlert(alert);
}
}
private void sendAlert(Alert alert) {
// 发送到告警系统
// 例如:钉钉、邮件、短信等
}
}
}
常见问题
Q1: 如何处理大表的初始同步?
解答:
- 增加并行度:使用多个 server-id 增加并行读取
- 调整分片大小:根据表大小设置合适的 splitSize
- 限流读取:控制 fetchSize 避免 OOM
- 分批同步:按时间范围或 ID 范围分批
// 大表同步优化
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("large_db")
.tableList("large_db.huge_table")
.username("root")
.password("password")
.serverId("5400-5415") // 16个并行度
.splitSize(4096) // 减小分片大小
.fetchSize(500) // 限制每次拉取量
.connectTimeout(Duration.ofMinutes(5))
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
Q2: 如何保证数据一致性?
解答:
- 启用 Checkpoint:确保故障恢复
- 使用 Exactly-Once 语义:端到端一致性
- 处理重复数据:使用主键去重
- 监控延迟:确保及时同步
Q3: Schema 变更如何处理?
解答:
- 自动感知:启用
include.schema.changes - 兼容性处理:使用 JSON 等灵活格式
- 版本管理:维护 Schema 版本历史
- 平滑升级:使用蓝绿部署
Q4: 性能优化建议?
解答:
- 合理设置并行度:根据源库负载调整
- 批量处理:减少网络开销
- 内存优化:调整缓冲区大小
- 监控指标:及时发现瓶颈
Q5: 如何处理删除操作?
解答:
// 处理删除事件
public class DeleteHandler implements DebeziumDeserializationSchema<String> {
@Override
public void deserialize(SourceRecord record, Collector<String> out) {
Envelope.Operation op = Envelope.operationFor(record);
if (op == Envelope.Operation.DELETE) {
// 处理删除
Struct before = value.getStruct("before");
String id = before.getString("id");
// 方案1:标记删除
JSONObject deleteEvent = new JSONObject();
deleteEvent.put("id", id);
deleteEvent.put("_deleted", true);
deleteEvent.put("_delete_time", System.currentTimeMillis());
out.collect(deleteEvent.toJSONString());
// 方案2:发送到删除主题
// sendToDeleteTopic(id);
}
}
}
相关文章
Flink 基础
- Flink 架构与工作流程 - 理解 Flink 核心架构
- Flink DataStream API - DataStream 编程基础
- Flink Table & SQL API - SQL 方式使用 CDC
高级特性
- Checkpoint & Savepoint - 容错机制详解
- Flink 性能优化 - 性能调优指南
- Flink 监控 - 监控体系建设
集成方案
- Flink + Kafka 集成 - 消息队列集成
- 实时数仓实战 - 完整案例
相关技术
- MySQL - MySQL 数据库
- PostgreSQL - PostgreSQL 数据库
- Kafka - 消息中间件
- Elasticsearch - 搜索引擎
总结
Flink CDC 提供了强大的实时数据同步能力,是构建实时数据管道的理想选择:
🎯 核心优势
- 实时性高:毫秒级延迟的数据同步
- 一致性强:Exactly-Once 语义保证
- 扩展性好:支持大规模分布式部署
- 使用简单:SQL 和 DataStream API 双支持
✅ 最佳实践
- 合理配置数据库参数(binlog、wal 等)
- 根据数据量选择合适的并行度
- 使用 Checkpoint 保证容错性
- 建立完善的监控告警机制
- 定期进行性能调优
🚀 适用场景
- 实时数据同步和备份
- 实时数仓和数据湖构建
- 缓存同步和搜索索引更新
- 微服务间数据共享
- 异构数据库迁移
🚫 注意事项:
- 确保源数据库开启必要的日志功能
- 监控源库压力,避免影响业务
- 合理设置资源,防止 OOM
- 做好 Schema 变更的兼容处理
通过 Flink CDC,可以轻松构建企业级的实时数据同步平台,为数据驱动的业务决策提供有力支撑。