Flink CDC 实时数据同步指南
在企业数据平台中,数据同步是最基础也是最关键的能力之一。传统的批量 ETL 虽然稳定,但无法满足实时决策的需求。Flink CDC(Change Data Capture)让你可以实时捕获数据库的变更事件,将其流式传输到 coomia-dip 的数据层,让业务数据库中的每一行变更都在毫秒级别反映到分析平台中。
“系列:S12 开发者教程 · 第 15 篇 | 难度:中级 | 阅读时间:15 分钟
Flink CDC 实时数据同步指南
#引言
在企业数据平台中,数据同步是最基础也是最关键的能力之一。传统的批量 ETL 虽然稳定,但无法满足实时决策的需求。Flink CDC(Change Data Capture)让你可以实时捕获数据库的变更事件,将其流式传输到 coomia-dip 的数据层,让业务数据库中的每一行变更都在毫秒级别反映到分析平台中。
本教程将带你理解 CDC 的原理,配置 MySQL binlog,部署 Flink 集群,并实现一个从 MySQL 到 Apache Doris 的实时同步管道。
#1. CDC 基础概念
#1.1 什么是 Change Data Capture
Change Data Capture 是一种通过监听数据库事务日志来捕获数据变更的技术。与定时全量拉取相比,CDC 有三个核心优势:
- 实时性:变更发生后毫秒级捕获,而非等待下一个调度周期
- 低负载:只传输变更数据,不需要全表扫描
- 完整性:捕获所有变更(包括中间状态),而非只有最终快照
CDC 的工作流程可以简化为:源数据库的事务日志被 Flink CDC Connector 持续消费,解析为结构化的变更事件(Changelog Stream),经过 Flink 的流处理进行清洗、转换和路由,最终写入目标系统(在 coomia-dip 中是 Apache Doris)。
#1.2 Flink CDC 3.x Pipeline 架构
Flink CDC 3.x 引入了声明式 Pipeline 架构,大幅简化了同步任务的配置。你只需要一个 YAML 文件即可描述整个同步拓扑:
source:
type: mysql
hostname: mysql-source
port: 3306
username: cdc_user
password: ${MYSQL_CDC_PASSWORD}
tables: erp_db.orders, erp_db.order_items, erp_db.customers
sink:
type: doris
fenodes: doris-fe:8030
username: root
password: ${DORIS_PASSWORD}
pipeline:
name: erp-to-doris-sync
parallelism: 4
schema.change.behavior: evolve
#1.3 coomia-dip 中的 CDC 定位
在 coomia-dip 的八平面架构中,CDC 属于 Data Layer(Data Layer)的核心能力。数据从业务数据库通过 CDC 进入 Doris ODS 层,然后由 Ontology Mapper 自动映射为 Ontology 对象,最终通过 OQL 引擎提供统一查询。
#2. 环境准备
#2.1 启用 MySQL Binlog
CDC 依赖 MySQL 的 binlog(二进制日志)。必须使用 ROW 格式以获取完整的字段级变更信息:
[mysqld]
server-id = 1
log_bin = mysql-bin
binlog_format = ROW
binlog_row_image = FULL
expire_logs_days = 3
gtid_mode = ON
enforce_gtid_consistency = ON
创建 CDC 专用数据库用户:
CREATE USER 'cdc_user'@'%' IDENTIFIED BY 'cdc_secure_password';
GRANT SELECT, RELOAD, SHOW DATABASES,
REPLICATION SLAVE, REPLICATION CLIENT
ON *.* TO 'cdc_user'@'%';
FLUSH PRIVILEGES;
验证配置:
SHOW VARIABLES LIKE 'log_bin'; -- ON
SHOW VARIABLES LIKE 'binlog_format'; -- ROW
SHOW VARIABLES LIKE 'binlog_row_image'; -- FULL
SHOW MASTER STATUS;
#2.2 部署 Flink 集群
使用 Docker Compose 快速部署 Flink Session 集群:
version: "3.8"
services:
jobmanager:
image: flink:1.18-java11
ports: ["8081:8081"]
command: jobmanager
environment:
FLINK_PROPERTIES: |
jobmanager.rpc.address: jobmanager
state.checkpoints.dir: file:///opt/flink/checkpoints
execution.checkpointing.interval: 60000
execution.checkpointing.mode: EXACTLY_ONCE
taskmanager:
image: flink:1.18-java11
depends_on: [jobmanager]
command: taskmanager
deploy:
replicas: 2
environment:
FLINK_PROPERTIES: |
jobmanager.rpc.address: jobmanager
taskmanager.numberOfTaskSlots: 4
taskmanager.memory.process.size: 4096m
#2.3 安装 CDC Connector
FLINK_LIB=/opt/flink/lib
wget -P $FLINK_LIB \
https://repo1.maven.org/maven2/com/ververica/flink-sql-connector-mysql-cdc/3.0.1/flink-sql-connector-mysql-cdc-3.0.1.jar
wget -P $FLINK_LIB \
https://repo1.maven.org/maven2/org/apache/doris/flink-doris-connector-1.18/1.6.1/flink-doris-connector-1.18-1.6.1.jar
#3. 实现同步管道
#3.1 Flink SQL 方式(推荐入门)
Flink SQL 是最快捷的 CDC 上手方式,不需要编写 Java 代码:
SET 'execution.runtime-mode' = 'streaming';
SET 'execution.checkpointing.interval' = '60s';
CREATE TABLE orders_source (
order_id BIGINT,
customer_id BIGINT,
product_name STRING,
quantity INT,
unit_price DECIMAL(10, 2),
total_amount DECIMAL(12, 2),
status STRING,
created_at TIMESTAMP(3),
updated_at TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql-source',
'port' = '3306',
'username' = 'cdc_user',
'password' = 'cdc_secure_password',
'database-name' = 'erp_db',
'table-name' = 'orders',
'scan.incremental.snapshot.enabled' = 'true',
'scan.incremental.snapshot.chunk.size' = '8096'
);
CREATE TABLE orders_sink (
order_id BIGINT,
customer_id BIGINT,
product_name STRING,
quantity INT,
unit_price DECIMAL(10, 2),
total_amount DECIMAL(12, 2),
status STRING,
created_at TIMESTAMP(3),
updated_at TIMESTAMP(3),
PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
'connector' = 'doris',
'fenodes' = 'doris-fe:8030',
'table.identifier' = 'ods.orders',
'username' = 'root',
'password' = '',
'sink.properties.format' = 'json',
'sink.properties.read_json_by_line' = 'true',
'sink.enable-delete' = 'true',
'sink.label-prefix' = 'flink_cdc_orders'
);
INSERT INTO orders_sink SELECT * FROM orders_source;
#3.2 DataStream API(灵活定制)
对于需要复杂转换逻辑的场景:
package com.onto.pipeline.cdc;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.cdc.connectors.mysql.source.MySqlSource;
import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema;
public class ERPSyncPipeline {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60_000, CheckpointingMode.EXACTLY_ONCE);
env.setParallelism(4);
MySqlSource<String> source = MySqlSource.<String>builder()
.hostname("mysql-source")
.port(3306)
.databaseList("erp_db")
.tableList("erp_db.orders", "erp_db.order_items", "erp_db.customers")
.username("cdc_user")
.password("cdc_secure_password")
.deserializer(new JsonDebeziumDeserializationSchema())
.includeSchemaChanges(true)
.splitSize(8096)
.build();
env.fromSource(source, WatermarkStrategy.noWatermarks(), "MySQL CDC")
.map(new CDCEventTransformer())
.process(new TableRoutingFunction())
.sinkTo(DorisSinkBuilder.build());
env.execute("ERP to coomia-dip CDC Pipeline");
}
}
转换函数示例:
public class CDCEventTransformer implements MapFunction<String, CDCRecord> {
@Override
public CDCRecord map(String json) throws Exception {
JSONObject event = JSON.parseObject(json);
String op = event.getString("op");
JSONObject after = event.getJSONObject("after");
JSONObject before = event.getJSONObject("before");
JSONObject src = event.getJSONObject("source");
CDCRecord record = new CDCRecord();
record.setDatabase(src.getString("db"));
record.setTable(src.getString("table"));
record.setOperation(op);
record.setData(after != null ? after : before);
record.setTimestamp(System.currentTimeMillis());
if ("orders".equals(record.getTable()) && after != null) {
double amount = after.getDoubleValue("total_amount");
after.put("amount_with_tax", amount * 1.13);
after.put("currency", "CNY");
}
return record;
}
}
#4. 与 coomia-dip Ontology 集成
#4.1 注册 CDC 数据源
from ontology_sdk import OntoPlatform
platform = OntoPlatform(base_url="http://localhost:8080", token="admin-token")
source = platform.data_sources.register(
name="erp-mysql-cdc",
source_type="FLINK_CDC",
config={
"connector": "mysql-cdc",
"hostname": "mysql-source",
"port": 3306,
"database": "erp_db",
"tables": ["orders", "order_items", "customers"],
"sync_mode": "INCREMENTAL",
"checkpoint_interval_ms": 60000,
},
target={"type": "DORIS", "database": "ods"},
)
#4.2 自动映射到 Ontology 对象
platform.ontology.create_mapping(
source_table="ods.orders",
object_type="Order",
field_mappings={
"order_id": {"property": "orderId", "type": "STRING", "primary_key": True},
"customer_id": {"property": "customerId", "type": "STRING"},
"product_name": {"property": "productName", "type": "STRING"},
"total_amount": {"property": "totalAmount", "type": "DOUBLE"},
"status": {"property": "status", "type": "STRING"},
"created_at": {"property": "createdAt", "type": "TIMESTAMP"},
},
sync_mode="REAL_TIME",
link_rules=[{
"link_type": "orderedBy",
"target_object_type": "Customer",
"source_field": "customer_id",
"target_field": "customerId",
}],
)
#4.3 OQL 实时查询
result = platform.oql.execute("""
SELECT o.orderId, o.productName, o.totalAmount, c.name AS customerName
FROM Order o
JOIN o.orderedBy c
WHERE o.totalAmount > 10000
AND o.createdAt > NOW() - INTERVAL 1 HOUR
ORDER BY o.totalAmount DESC
LIMIT 20
""")
for row in result.rows:
print(f"Order {row['orderId']}: {row['customerName']} - {row['totalAmount']:.2f}")
#5. 监控与运维
#5.1 关键指标
| 指标 | 含义 | 告警阈值 |
|---|---|---|
numRecordsInPerSecond | 每秒输入记录数 | 突变 > 200% 或降为 0 |
numRecordsOutPerSecond | 每秒输出记录数 | 与输入差距 > 20% |
currentFetchEventTimeLag | CDC 事件延迟 | > 60 秒 |
lastCheckpointDuration | Checkpoint 耗时 | > 30 秒 |
numberOfFailedCheckpoints | 失败次数 | > 0 |
#5.2 常见问题排查
| 问题 | 原因 | 解决方案 |
|---|---|---|
| 任务启动失败 | binlog 已过期 | 增大 expire_logs_days,全量重同步 |
| 全量阶段 OOM | 表太大 | 减小 chunk.size,增加 TM 内存 |
| Doris 写入超时 | Stream Load 慢 | 增大超时参数,减小批次 |
| Schema 变更中断 | DDL 处理异常 | 启用 schema.change.behavior: evolve |
| 数据延迟增大 | Sink 吞吐不足 | 增加并行度 |
| 数据重复 | Checkpoint 不一致 | 启用 Exactly-Once + 2PC |
#5.3 灾难恢复
# 创建 Savepoint
curl -X POST http://flink-jobmanager:8081/jobs/{JOB_ID}/savepoints \
-d '{"cancel-job": false, "target-directory": "file:///opt/flink/savepoints"}'
# 从 Savepoint 恢复
./bin/flink run \
-s file:///opt/flink/savepoints/savepoint-xxxxxx \
-c com.onto.pipeline.cdc.ERPSyncPipeline \
erp-sync-pipeline.jar
#6. 生产环境最佳实践
#6.1 容量规划
| 源表规模 | 变更 TPS | 推荐并行度 | TM 内存 |
|---|---|---|---|
| < 100 万 | < 100 | 2 | 2 GB |
| 100-1000 万 | 100-1000 | 4 | 4 GB |
| 1000 万-1 亿 | 1000-5000 | 8 | 8 GB |
| > 1 亿 | > 5000 | 16+ | 16 GB+ |
#6.2 上线检查清单
- MySQL binlog 格式为 ROW,
binlog_row_image为 FULL - binlog 保留至少 3 天
- CDC 用户权限最小化
- Flink checkpoint 使用持久化存储
- 监控告警已配置
- Doris 目标库 Schema 已创建
- 网络连通性已验证
- Savepoint 恢复流程已演练
#总结
本教程完整介绍了使用 Flink CDC 实现实时数据同步的全流程:CDC 原理、环境配置、Flink SQL 和 DataStream API 两种实现方式、Ontology 自动映射、运维监控和生产最佳实践。Flink CDC 是 coomia-dip 数据层的"血液循环系统",让企业各业务系统的数据实时汇聚到统一的 Ontology 模型中。
下一篇:[S12-16] Temporal 工作流编排指南 上一篇:[S12-14] gRPC 自定义服务开发指南