易君召
易君召
发布于 2026-07-19 / 3 阅读
0
0

Apache Flink CDC 数据同步传输协议简要说明

Apache Flink CDC 的数据同步传输协议需要分三层来看,整体以 TCP 协议为核心,并非 HTTP:

一、数据采集层(Source 端:CDC 读取数据库)

Flink CDC 内置 Debezium 引擎,直接与源数据库建立长连接读取变更日志,底层均为 TCP 协议,但不同数据库使用各自的应用层协议:

数据库

应用层协议

底层传输

说明

MySQL

MySQL Binlog Dump 协议

TCP(默认 3306 端口)

发送 COM_BINLOG_DUMP 命令订阅 binlog,本质是 MySQL 主从复制协议的客户端实现

PostgreSQL

逻辑复制协议(START_REPLICATION)

TCP(默认 5432 端口)

以 "伪备库" 身份连接,通过复制槽流式接收 WAL 日志

Oracle

JDBC + LogMiner / XStream

TCP(默认 1521 端口)

通过 JDBC 连接,调用数据库日志挖掘接口获取变更

SQL Server

CDC 函数调用(JDBC)

TCP

轮询 SQL Server 内置的 CDC 变更表

核心结论:CDC 采集阶段不使用 HTTP,全部是数据库原生协议的 TCP 长连接,实时性高、开销低。

CDC 数据进入 Flink 后,在算子之间的传输分为两种情况:

  1. 跨 TaskManager 传输

    • 基于 Netty 框架实现,底层为 TCP 协议

    • 使用 Flink 自定义的 NettyMessage 二进制消息格式

    • 支持零拷贝、缓冲复用等高性能优化CSDN博...

  2. 同 TaskManager 内传输

    • 若算子被链化(Operator Chain),直接通过方法调用传递数据,不经过网络

    • 无序列化和网络开销,性能最优

此外,JobManager 与 TaskManager 之间的控制信令(如任务提交、心跳)基于 Pekko/Akka RPC,底层同样是 TCP。

三、数据写出层(Sink 端:写入下游系统)

这一层协议取决于具体的下游存储,差异较大:

  • 写入 Kafka / Pulsar:使用消息队列原生协议,底层 TCP

  • 写入 Doris / StarRocks:通过 HTTP Stream Load 批量导入Apach...

  • 写入 Elasticsearch / OpenSearch:HTTP REST API

  • 写入 JDBC 数据库:JDBC 协议,底层 TCP

  • 写入 HDFS / OSS / S3:对应分布式存储协议

总结一句话

Flink CDC 采集数据走数据库原生 TCP 协议,Flink 内部数据流转走 Netty TCP;只有写入部分下游系统(如 Doris、ES)时才会用到 HTTP。 整个链路中 TCP 是绝对主流,HTTP 仅出现在特定 Sink 场景。


评论