Flink CDC Kafka Pipeline Connector 全指南:构建 MySQL 到 Kafka 的实时数据管道

📅 发布时间:2026/9/17 6:17:53
Flink CDC Kafka Pipeline Connector 全指南:构建 MySQL 到 Kafka 的实时数据管道
Flink CDC Kafka Pipeline Connector 全指南构建 MySQL 到 Kafka 的实时数据管道【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdcKafka Pipeline Connector 是 Flink CDC Pipeline 体系中面向 Apache Kafka 的内置 Data Sink用于将上游数据库如 MySQL的变更数据持续写入 Kafka Topic实现数据库 → Kafka的实时数据同步。本文基于当前仓库中 Kafka Pipeline Connector 的官方文档kafka.md并结合其源码实现完整讲解 Pipeline 定义、全部可配置项、输出消息格式与数据类型映射帮助你从零构建一个可投入生产使用的 MySQL → Kafka 同步管道并深入理解其底层工作原理。Kafka Pipeline Connector 能做什么Kafka Pipeline Connector 作为 Pipeline 的Data Sink数据汇使用核心能力是数据同步Data synchronization它将上游数据源产生的Event数据变更事件与 Schema 变更事件序列化后写入 Kafka让下游消费者能够以统一的 Kafka 消息流获取数据库的实时变更。在架构上Pipeline 中的数据流是Source → 路由/转换 → SinkKafka Connector 位于整条链路的最末端。它基于 Flink 官方 Kafka Connector 的KafkaSink构建见 KafkaDataSink.java因此天然继承了 Kafka 生产端的可靠性机制。关于 Pipeline 的整体概念可参考>source: type: mysql name: MySQL Source hostname: 127.0.0.1 port: 3306 username: admin password: pass tables: adb.\.*, bdb.user_table_[0-9], [app|web].order_\.* server-id: 5401-5404 sink: type: kafka name: Kafka Sink properties.bootstrap.servers: PLAINTEXT://localhost:62510 pipeline: name: MySQL to Kafka Pipeline parallelism: 2要点解读source 部分tables支持正则与通配符组合server-id: 5401-5404为多并行度场景下的 binlog 客户端 ID 区间与pipeline.parallelism: 2相匹配。sink 部分type: kafka声明使用本 Connectorproperties.bootstrap.servers为必填项指定 Kafka 集群的 broker 地址列表。pipeline 部分parallelism: 2设置整条管道含 Source 与 Sink的并行度。将上述内容保存为 YAML 文件后即可通过 Flink CDC 提供的 CLI 工具提交运行完整的环境准备与提交步骤可参考仓库中的快速入门教程 mysql-to-kafka.mdFlink 2.2 版本或 quickstart-for-1.20 下的 mysql-to-kafka.md。Pipeline Connector 配置项全解析以下是 Kafka Pipeline Connector 支持的全部配置项信息以官方文档为准并补充了 KafkaDataSinkOptions.java 源码中的默认值与类型定义OptionRequiredDefaultTypeDescriptiontyperequired(none)String指定使用的 Connector 类型此处必须为kafka。nameoptional(none)StringSink 的名称。partition.strategyoptional(none)String定义记录发送到 Kafka 分区的策略。可选值all-to-zero将所有记录发送到 0 号分区、hash-by-key按主键哈希将记录分布到各分区默认值为all-to-zero。key.formatoptional(none)String定义 key 数据的编码格式标识。可选值csv、json默认值为json。value.formatoptional(none)String用于序列化 Kafka 消息 value 部分的格式。可选值debezium-json、canal-json默认值为debezium-json目前不支持自定义格式。properties.bootstrap.serversrequired(none)String用于与 Kafka 集群建立初始连接的主机/端口对列表。topicoptional(none)String若配置了该参数所有事件都会被发送到该 Topic。sink.add-tableId-to-header-enabledoptional(none)Boolean若为 true则为每条 Kafka 记录添加 key 为namespace、schemaName、tableName的 Header默认值为 false。properties.*optional(none)String透传给 Kafka 生产端的配置项即 Kafka 生产者相关配置。sink.custom-headeroptional(none)String为每条 Kafka 记录添加的自定义 Header。多个 Header 用,分隔key 与 value 用:分隔例如key1:value1,key2:value2。sink.tableId-to-topic.mappingoptional(none)String为每个上游表tableId自定义其对应的下游 Kafka Topic 映射。多条映射用;分隔上游 tableId 与下游 Topic 用:分隔例如mydb.mytable1:topic1;mydb.mytable2:topic2。debezium-json.include-schema.enabledoptionalfalseBoolean若配置该参数每条 debezium 记录将包含 Debezium Schema 信息。仅在value.format为debezium-json时支持。各配置项的源码级解读从 KafkaDataSinkFactory.java 的createDataSink方法可以看到这些配置项是如何被消费的properties.*透传机制工厂会遍历所有以properties.前缀开头的配置项去掉前缀后原样放入 Kafka Producer 的Properties中对应PROPERTIES_PREFIX常量与 KafkaDataSinkFactory.java 的逻辑。这意味着你可以通过properties.acks、properties.linger.ms、properties.compression.type等任意标准 Kafka Producer 参数微调写入行为同时properties.bootstrap.servers会被单独取出用于构建KafkaSinkBuilder。分区策略实现partition.strategy的取值在 PartitionStrategy.java 中定义为枚举all-to-zero与hash-by-key。序列化阶段PipelineKafkaRecordSerializationSchema.java会将all-to-zero映射为固定分区 0而hash-by-key则交由 Kafka 生产端按 key 哈希自动路由。时区联动value.format序列化时会读取 Pipeline 级的pipeline.local-time-zone配置作为时区若未设置则取系统默认时区该时区会传入 Debezium/Canal JSON 序列化器中影响时间戳类字段的格式化输出。隐藏的可靠性参数源码中还定义了一个文档表格之外的参数sink.delivery-guarantee枚举类型默认AT_LEAST_ONCE可选EXACTLY_ONCE等见 KafkaDataSinkOptions.java它直接对应底层KafkaSinkBuilder.setDeliveryGuarantee(...)用于控制写入 Kafka 的交付保证级别。使用注意事项Topic 命名规则Kafka 写入的 Topic 默认为 TableId 的字符串表示namespace.schemaName.tableName。如果想改变 Topic 名称可以使用 Pipeline 的 route 功能进行改写相关用法见 route.md。Topic 自动创建如果写入的 Kafka Topic 不存在Connector 会自动创建前提是 Kafka broker 允许自动创建 Topic。Schema 变更事件处理从源码看PipelineKafkaRecordSerializationSchema.javaSchemaChangeEvent建表、加列等 DDL 事件不会作为独立消息写入 Kafka而是被用于更新内部的表结构缓存以驱动后续数据行的序列化。Topic 推断优先级对应 PipelineKafkaRecordSerializationSchema.java配置了topic统一 Topic→ 命中sink.tableId-to-topic.mapping的映射规则 → 兜底使用 TableId 字符串。其中映射规则通过Selectors支持正则匹配表名解析解析逻辑在 KafkaSinkUtils.java 中实现并且结果会按 TableId 做缓存以提升性能。自定义 Headersink.add-tableId-to-header-enabledtrue时会在每条记录上附加namespace、schemaName、tableName三个 Header空值以空字符串代替sink.custom-header则用于附加任意自定义键值对两者可同时生效。输出消息格式Output FormatKafka 消息的 value 部分根据value.format的不同而不同。key 部分由key.format决定json或csv。debezium-json 格式该格式遵循 Debezium 消息结构包含before、after、op、source等元素但注意source中不包含ts_ms字段。op的取值在 DebeziumJsonSerializationSchema.java 中定义为cinsert、uupdate、ddelete。输出示例{ before: null, after: { col1: 1, col2: 1 }, op: c, source: { db: default_namespace, table: table1 } }当debezium-json.include-schema.enabled为true时输出会额外包含schema与payload两层结构即 Debezium Envelope 完整格式{ schema:{ type:struct, fields:[ { type:struct, fields:[ { type:string, optional:true, field:col1 }, { type:string, optional:true, field:col2 } ], optional:true, field:before }, { type:struct, fields:[ { type:string, optional:true, field:col1 }, { type:string, optional:true, field:col2 } ], optional:true, field:after } ], optional:false }, payload:{ before: null, after: { col1: 1, col2: 1 }, op: c, source: { db: default_namespace, table: table1 } } }canal-json 格式该格式遵循 Canal 消息结构包含old、data、type、database、table、pkNames等元素但不包含ts字段。type的取值在 CanalJsonSerializationSchema.java 中定义为INSERT、UPDATE、DELETE。输出示例{ old: null, data: [ { col1: 1, col2: 1 } ], type: INSERT, database: default_schema, table: table1, pkNames: [ col1 ] }两种格式的序列化实现分别位于flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-kafka/src/main/java/org/apache/flink/cdc/connectors/kafka/json/debezium/与.../json/canal/目录下选择逻辑由 JsonSerializationType.java 枚举驱动。需要说明的是Debezium/Canal 的序列化均按表维度动态维护 schema 状态——当收到CreateTableEvent时构建序列化器收到其他SchemaChangeEvent如 ALTER TABLE时基于已缓存的 schema 增量更新从而保证字段演化的正确性。因此在编写下游消费者时建议同时消费 schema 变更信号或使用兼容 JSON Schema 演化的消费端工具。数据类型映射Data Type Mapping下表给出 CDC 类型在 Kafka JSON 消息中的映射关系。其中Literal type字面类型定义了数据的物理存储格式对应 Debezium schema 的type字段Semantic type语义类型定义了数据的逻辑含义对应 Debezium schema 的name字段CDC typeJSON typeLiteral typeSemantic typeNOTETINYINTTINYINTINT16SMALLINTSMALLINTINT16INTINTINT32BIGINTBIGINTINT64FLOATFLOATFLOATDOUBLEDOUBLEDOUBLEDECIMAL(p, s)DECIMAL(p, s)BYTESorg.apache.kafka.connect.data.DecimalBOOLEANBOOLEANBOOLEANDATEDATEINT32io.debezium.time.DateTIMESTAMP(p)TIMESTAMP(p)INT64p 3 时io.debezium.time.Timestampp 3 时io.debezium.time.MicroTimestampTIMESTAMP_LTZTIMESTAMP_LTZSTRINGio.debezium.time.ZonedTimestampCHAR(n)CHAR(n)STRINGVARCHAR(n)VARCHAR(n)STRINGARRAYARRAYARRAY序列化为 JSON 数组元素类型按本表递归转换。MAPMAPMAP序列化为 JSON 对象key 与 value 类型按本表递归转换。ROWROWSTRUCT序列化为 JSON 对象字段类型按本表递归转换。该映射关系与序列化器中的实际转换逻辑一致例如 DebeziumJsonSerializationSchema.java 中引入了io.debezium.data.Bits、io.debezium.time.Date/MicroTime/MicroTimestamp/Timestamp/ZonedTimestamp以及 Kafka Connect 的Decimal等类型来构造 Debezium Schema与上表一一对应。小结Kafka Pipeline Connector 以极简的 YAML 配置即可完成数据库 → Kafka的实时数据同步同时通过partition.strategy、key.format、value.format、topic、sink.tableId-to-topic.mapping、sink.custom-header、debezium-json.include-schema.enabled等参数提供了分区路由、消息格式、Topic 定制、Header 附加等灵活的调优手段。结合 KafkaDataSinkOptions.java 与序列化器源码你可以进一步确认每条配置的真实语义为下游 Flink SQL、数据湖或消息消费场景构建可靠的变更数据流。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考