易君召
易君召
发布于 2026-07-14 / 12 阅读
0
0

Apache Camel + Flink CDC 技术结合方案选型

一、两种主流融合架构

适用:数据库 binlog 实时捕获,由 Camel 负责异构目的地分发(Kafka、MySQL、MinIO、API、FTP 等) 流程: MySQL/PG Binlog → Flink CDC Source → Flink Stream → 输出到 Camel Producer(Camel Flink Connector)

适用:小数据量、无需复杂流计算,直接用 Camel 内置 Debezium 组件捕获 binlog,省去部署 Flink 集群 Camel Debezium Component → 路由转换 → Kafka/DB/文件

Camel 消费 FTP/SFTP/HTTP/MQ 数据,送入 Flink 做复杂聚合、窗口计算

下面分两套方案给出落地代码与原理。

二、方案一Flink CDC + Camel Flink Connector(生产主流,大数据实时同步)

1. 核心原理

  1. Flink CDC Source 捕获数据库增量 binlog(全量 + 增量)

  2. 自定义 Flink Sink:使用 camel-flink 连接器,把 Flink DataStream 消息交给 Camel Producer

  3. Camel 路由统一处理:数据清洗、格式转换、路由分发到任意异构终端(Kafka、ES、MySQL、MinIO、第三方 HTTP 接口)

<!-- 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>
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. 优势

  1. Flink CDC 专业稳定:支持全量快照、断点续传、分库分表、高并发读取 binlog

  2. Camel 屏蔽异构输出端:一套路由同时输出 MQ、数据库、对象存储、HTTP 接口

  3. 解耦:流计算(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 + Camel Sink

Camel Debezium 直连 CDC

集群依赖

需要部署 Flink 集群

仅 SpringBoot 微服务

数据吞吐量

高,支持分库分表海量数据

中低,单进程消费

流计算能力

支持窗口、join、复杂状态计算

无流计算,仅简单转换过滤

断点续传

Checkpoint 完整 Exactly-Once

文件偏移量,仅 At-Least-Once

适用场景

企业级大数据实时数据管道、数据中台、可信数据空间

小型系统、轻量同步、低成本快速落地

六、典型业务落地场景

  1. 数据中台同步:Flink CDC 捕获多库 binlog → Camel 分发 Kafka、ES、数仓、MinIO 归档

  2. 可信数据空间连接器:CDC 增量数据经 Camel 脱敏、格式标准化后推送到数据共享总线

  3. 异构系统实时互通:数据库变更自动同步至第三方 HTTP 业务平台、FTP 对账文件

  4. IoT 工业场景:OPC UA 设备数据(Camel)+ 数据库 CDC(Flink)双流进入 Flink 做联动计算

七、核心注意点

  1. 一致性:Flink Checkpoint + Camel 事务生产者,实现端到端精准一次;

  2. 数据转换:统一使用 Jolt 组件做异构 JSON 对齐,适配多系统字段差异;

  3. 偏移存储:Flink 使用 Checkpoint,Camel Debezium 使用文件 / Kafka 存储 offset,防止重复消费;

  4. 异常容错:Camel 路由配置重试、死信队列,CDC 异常消息单独落盘不阻塞整条同步链路。


原文链接 https://www.yijunzhao.cn/archives/apache-camel-flink-cdc-integration-comparison

欢迎访问 小易撩挨踢

https://www.yijunzhao.cn/


评论