一、两种主流融合架构
架构 1:Flink CDC 做数据采集 → Camel 做数据分发 / 转换(最常用)
适用:数据库 binlog 实时捕获,由 Camel 负责异构目的地分发(Kafka、MySQL、MinIO、API、FTP 等) 流程: MySQL/PG Binlog → Flink CDC Source → Flink Stream → 输出到 Camel Producer(Camel Flink Connector)
架构 2:Camel Debezium 直连数据库 CDC(轻量无 Flink 集群)
适用:小数据量、无需复杂流计算,直接用 Camel 内置 Debezium 组件捕获 binlog,省去部署 Flink 集群 Camel Debezium Component → 路由转换 → Kafka/DB/文件
架构 3:Camel 作为数据源入 Flink
Camel 消费 FTP/SFTP/HTTP/MQ 数据,送入 Flink 做复杂聚合、窗口计算
下面分两套方案给出落地代码与原理。
二、方案一Flink CDC + Camel Flink Connector(生产主流,大数据实时同步)
1. 核心原理
Flink CDC Source 捕获数据库增量 binlog(全量 + 增量)
自定义 Flink Sink:使用
camel-flink连接器,把 Flink DataStream 消息交给 Camel ProducerCamel 路由统一处理:数据清洗、格式转换、路由分发到任意异构终端(Kafka、ES、MySQL、MinIO、第三方 HTTP 接口)
2. 依赖(Flink 1.18 + Camel 4.x)
<!-- Flink CDC MySQL -->
<dependency>
<groupId>com.ververica</groupId>
<artifactId>flink-connector-mysql-cdc</artifactId>
<version>2.4.0</version>
</dependency>
<!-- Flink Stream Core -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>1.18.0</version>
</dependency>
<!-- Camel Flink 连接器 -->
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-flink</artifactId>
<version>4.4.0</version>
</dependency>
<!-- 目标端组件,示例Kafka -->
<dependency>
<groupId>org.apache.camel</groupId>
<artifactId>camel-kafka</artifactId>
<version>4.4.0</version>
</dependency>
3. 完整代码示例:MySQL CDC → Flink → Camel → Kafka
import org.apache.camel.flink.CamelSink;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class FlinkCamelCdcJob {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(3000);
// 1. 构建Flink CDC MySQL Source
MySqlSource<String> mysqlSource = MySqlSource.<String>builder()
.hostname("127.0.0.1")
.port(3306)
.databaseList("test_db")
.tableList("test_db.user")
.username("root")
.password("123456")
.deserializer(new JsonDebeziumDeserializationSchema())
.build();
// 2. 读取binlog数据流
env.fromSource(mysqlSource, org.apache.flink.streaming.api.environment.StreamSourceConfig.newBuilder().build(), "mysql-cdc")
// 3. Camel Sink:把binlog消息交给Camel路由
.addSink(CamelSink.forUri("kafka:user_cdc_topic?brokers=127.0.0.1:9092")
.withCamelContextConfigurer(context -> {
// 可注册自定义转换器、数据清洗路由
context.getRegistry().bind("joltTransform", new JoltDataTransform());
})
.build());
env.execute("Flink-Camel-MySQL-CDC-Job");
}
}
4. 扩展:在 Camel 内部做复杂数据处理
通过 Camel 的路由能力在 Sink 内部完成转换,不用写 Flink 算子:
// 不直接写kafka,交给direct路由做清洗后转发
CamelSink.forUri("direct:cdcProcess")
.withCamelContextConfigurer(context -> {
context.addRoutes(new RouteBuilder() {
@Override
public void configure() throws Exception {
from("direct:cdcProcess")
// 1. Jolt JSON字段映射、脱敏
.transform().jolt("classpath:cdc/jolt/user-mapping.json")
// 2. 过滤删除事件
.filter(simple("${header.operation} != 'd'"))
// 3. 写入Kafka
.to("kafka:user_cdc_topic")
// 4. 同步写入MySQL备份表
.to("sql:INSERT INTO user_backup(id,name) VALUES(:#id,:#name)");
}
});
})
5. 优势
Flink CDC 专业稳定:支持全量快照、断点续传、分库分表、高并发读取 binlog
Camel 屏蔽异构输出端:一套路由同时输出 MQ、数据库、对象存储、HTTP 接口
解耦:流计算(Flink)和集成分发(Camel)职责分离
三、方案二轻量方案 Camel Debezium 内置 CDC(无 Flink 集群)
不需要部署 Flink,SpringBoot+Camel 直接捕获数据库 binlog,适合中小业务、轻量化同步场景。
1. 依赖
<dependency>
<groupId>org.apache.camel.springboot</groupId>
<artifactId>camel-debezium-mysql-spring-boot-starter</artifactId>
<version>4.4.0</version>
</dependency>
2. 配置 application.yml
camel:
component:
debezium-mysql:
database-hostname: 127.0.0.1
database-port: 3306
database-user: root
database-password: 123456
database-server-name: mysql-server-1
database-include-list: test_db
table-include-list: test_db.user
offset-storage: file
offset-file-filename: ./offset/offset.dat
3. Camel CDC 路由代码
@Component
public class CamelDebeziumRoute extends RouteBuilder {
@Override
public void configure() throws Exception {
// 监听MySQL binlog
from("debezium-mysql:mysql-server-1")
.log("捕获CDC变更数据: ${body}")
// JSON转换、清洗
.transform().jolt("classpath:cdc/jolt/user.json")
// 分发到Kafka
.to("kafka:user_topic")
// 同步写入ES
.to("elasticsearch:localhost:9200/user_index");
}
}
局限
不支持 Flink 强大的窗口、双流 join、复杂状态计算,仅适合单纯数据同步。
四、方案三Camel 作为 Flink 数据源(反向结合)
场景:外部异构源(FTP、SFTP、第三方 API、MQTT IoT 设备)数据先由 Camel 采集,推送至 Flink 做实时计算。 流程: FTP/SFTP/HTTP → Camel Route → Camel Flink Source → Flink流处理
示例路由:Camel 读取 FTP 文件,推送到 Flink 流
from("ftp://user@127.0.0.1/in?password=xxx")
.to("flink-stream:ftp-data-stream");
Flink 代码读取 flink-stream 数据源做窗口统计、聚合。
五、技术方案选型对比
六、典型业务落地场景
数据中台同步:Flink CDC 捕获多库 binlog → Camel 分发 Kafka、ES、数仓、MinIO 归档
可信数据空间连接器:CDC 增量数据经 Camel 脱敏、格式标准化后推送到数据共享总线
异构系统实时互通:数据库变更自动同步至第三方 HTTP 业务平台、FTP 对账文件
IoT 工业场景:OPC UA 设备数据(Camel)+ 数据库 CDC(Flink)双流进入 Flink 做联动计算
七、核心注意点
一致性:Flink Checkpoint + Camel 事务生产者,实现端到端精准一次;
数据转换:统一使用 Jolt 组件做异构 JSON 对齐,适配多系统字段差异;
偏移存储:Flink 使用 Checkpoint,Camel Debezium 使用文件 / Kafka 存储 offset,防止重复消费;
异常容错:Camel 路由配置重试、死信队列,CDC 异常消息单独落盘不阻塞整条同步链路。
原文链接
欢迎访问 小易撩挨踢