Kafka CDC Pipeline Connector 可以在 Flink CDC Pipeline 中作为 Source 消费 Kafka 消息,也可以作为 Sink 将数据库变更和 Schema 信息序列化后写入 Kafka。本文统一介绍两种角色的配置、格式、TableId、Topic 路由和使用方式。
类别 | Kafka Source | Kafka Sink |
|---|---|---|
支持类型 | Source | Sink |
作业类型 | Flink CDC Pipeline YAML | Flink CDC Pipeline YAML |
运行模式 | 流模式 | 流模式 |
数据方向 | Kafka → Flink CDC Pipeline | Flink CDC Pipeline → Kafka |
变更类型 | 根据消息格式解析 INSERT、UPDATE、DELETE,并根据字段变化推导 Schema | 将 INSERT、UPDATE、DELETE 序列化为 CDC JSON;Schema 输出取决于 Source 和格式配置 |
数据格式 | JSON、Debezium JSON、Canal JSON | Debezium JSON、Canal JSON |
组件形式 | 平台内置 | 平台内置 |
交付语义 | 启用 Checkpoint 后提供 Exactly-Once Source 处理能力 | At-Least-Once |
组件 | 支持版本 | 说明 |
|---|---|---|
火山引擎流式计算 Flink 版 |
| Source 和 Sink 均为平台内置 Connector,无需额外上传 JAR。 |
Flink CDC | 3.4.0 及以上 | 创建作业时请选择 3.4.0 或更高版本。 |
Kafka | 支持火山引擎 Kafka、自建 Kafka,以及兼容 Kafka 协议的服务。 |
topic 与 topic-pattern 必须二选一,不能同时配置。table-id.fixed.name,或同时配置 table-id.dynamic.from 和 table-id.dynamic.fields。table-id.dynamic.fields 支持 1~3 个逗号分隔的 JSON Path,目标字段值必须为字符串。latest-offset 时不会补读作业启动前的消息;使用 timestamp 或 specific-offsets 前应确认位点仍在 Kafka 保留范围内。topic 后,所有表的消息都会写入同一个 Topic,需要通过消息中的 TableId 信息区分来源。source: type: kafka name: Kafka Source properties.bootstrap.servers: <kafka_bootstrap_servers> properties.group.id: <consumer_group_id> topic: <source_topic> value.format: debezium-json table-id.dynamic.from: value table-id.dynamic.fields: source.db,source.table scan.startup.mode: group-offsets
sink: type: kafka name: Kafka Sink properties.bootstrap.servers: <kafka_bootstrap_servers> value.format: debezium-json topic: <target_topic>
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
|---|---|---|---|---|---|
| Pipeline Connector 类型。 | STRING | 是 | 无 | 固定为 |
| Source 名称。 | STRING | 否 | 无 | 用于作业拓扑和日志识别。 |
| Kafka Broker 接入地址。 | STRING | 是 | 无 | 多个地址使用英文逗号分隔。 |
| Kafka Consumer Group ID。 | STRING | 是 | 无 | 不同业务链路建议使用不同 Group ID。 |
| 订阅的 Topic。 | STRING / LIST | 条件必填 | 无 | 与 |
| 订阅 Topic 的正则表达式。 | STRING | 条件必填 | 无 | 与 |
| 新 Topic 或 Partition 的发现间隔。 | DURATION | 否 |
| 设置为 |
| Kafka Consumer 原生参数。 | STRING | 否 | 无 | 去除 |
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
|---|---|---|---|---|---|
| Kafka Key 解析格式。 | STRING | 否 |
| 支持 |
| Key 字段前缀。 | STRING | 否 | 无 | 用于避免 Key 与 Value 字段重名。 |
| Kafka Value 解析格式。 | STRING | 否 |
| 建议使用 |
| Value 字段前缀。 | STRING | 否 | 无 | 用于避免 Key 与 Value 字段重名。 |
| 为所有消息设置固定 TableId。 | STRING | 条件必填 | 无 | 与动态 TableId 配置二选一。 |
| 动态 TableId 的提取位置。 | STRING | 条件必填 | 无 | 支持 |
| 动态 TableId 的字段 Path。 | STRING | 条件必填 | 无 | 支持 1~3 个逗号分隔 Path。 |
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
|---|---|---|---|---|---|
| Kafka 消费起始位点。 | STRING | 否 |
| 支持 |
| 时间戳启动位点。 | LONG | 条件必填 | 无 |
|
| 指定 Partition Offset。 | STRING | 条件必填 | 无 | 示例: |
| 透传到下游的 Kafka Metadata。 | STRING | 否 | 无 | 支持 |
| Source 限速值。 | LONG | 否 | 无 | 单位为记录数。 |
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
|---|---|---|---|---|---|
| 是否跳过无法解析的字段和消息。 | BOOLEAN | 否 |
| 开启后错误消息可能被丢弃。 |
| 是否将原始类型统一推导为 STRING。 | BOOLEAN | 否 |
| — |
| 输入消息是否包含 Kafka Connect Schema。 | BOOLEAN | 否 |
| 应与上游消息结构保持一致。 |
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
|---|---|---|---|---|---|
| 是否跳过无法解析的字段和消息。 | BOOLEAN | 否 |
| 开启后错误消息可能被丢弃。 |
| 是否将原始类型统一推导为 STRING。 | BOOLEAN | 否 |
| — |
| Schema 推导策略。 | STRING | 否 |
| 支持 |
| 是否将 MySQL | BOOLEAN | 否 |
| 仅在 |
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
|---|---|---|---|---|---|
| Pipeline Connector 类型。 | STRING | 是 | 无 | 固定为 |
| Sink 名称。 | STRING | 否 | 无 | 用于作业拓扑和日志识别。 |
| Kafka Broker 接入地址。 | STRING | 是 | 无 | 多个地址使用英文逗号分隔。 |
| Kafka Value 序列化格式。 | STRING | 否 |
| 支持 |
| 固定写入的 Topic。 | STRING | 否 | 无 | 不配置时按 TableId 生成 Topic。 |
| 是否将 TableId 写入 Header。 | BOOLEAN | 否 |
| Header Key 为 |
| 自定义 Kafka Header。 | STRING | 否 | 无 | 示例: |
| Kafka Producer 原生参数。 | STRING | 否 | 无 | 去除 |
参数 | 说明 | 数据类型 | 是否必填 | 默认值 | 备注 |
|---|---|---|---|---|---|
| Kafka 分区策略。 | STRING | 否 |
|
|
| Kafka Key 序列化格式。 | STRING | 否 |
| 支持 |
| TableId 到 Topic 的映射。 | STRING | 否 | 无 | 多个映射用分号分隔;映射到的 Topic 须提前创建。 |
| Debezium JSON 是否携带 Schema。 | BOOLEAN | 否 |
| 仅适用于 |
单表 Topic 可以配置固定 TableId:
table-id.fixed.name: <schema_name>.<table_name>
动态 TableId 从 Kafka Key 或 Value 中提取:
table-id.dynamic.from: value table-id.dynamic.fields: source.db,source.table
字段数量决定生成方式:
字段数量 | 生成结果 |
|---|---|
1 | 将字段值解析为完整 TableId。 |
2 | 生成 |
3 | 生成 |
topic:所有消息写入固定 Topic。topic:默认根据上游 TableId 生成 Topic;对应 Topic 须提前创建。sink.tableId-to-topic.mapping 直接配置 TableId 到 Topic 的映射;映射到的 Topic 须提前创建。debezium-json 输出 before、after、op、source;canal-json 输出 old、data、type、database、table、pkNames。
Flink CDC 类型 | Kafka JSON 表示 |
|---|---|
TINYINT、SMALLINT、INT、BIGINT | JSON Number |
FLOAT、DOUBLE、DECIMAL | JSON Number 或对应 CDC 格式表示 |
BOOLEAN | JSON Boolean |
DATE、TIMESTAMP、TIMESTAMP_LTZ | CDC 格式定义的日期时间表示 |
CHAR、VARCHAR | JSON String |
ARRAY | JSON Array |
MAP、ROW | JSON Object |
以下示例从 Kafka 消费 Debezium JSON,根据 source.db 和 source.table 生成 TableId,并写入 Values Sink 验证。
source: type: kafka name: Kafka Source properties.bootstrap.servers: <kafka_bootstrap_servers> properties.group.id: <consumer_group_id> topic: <source_topic> value.format: debezium-json table-id.dynamic.from: value table-id.dynamic.fields: source.db,source.table scan.startup.mode: group-offsets metadata.list: topic,partition,offset,timestamp sink: type: values name: Values Sink pipeline: name: Kafka CDC Source Pipeline parallelism: 2
预期结果。 Source 从 Consumer Group 已提交位点读取消息,解析 CDC Row Change 和 TableId,并将 Kafka Metadata 传递给下游。
以下示例使用占位 Source 表示任意受支持的数据库 CDC Source,实际运行时替换为对应 Source 参数。
source: type: <database_cdc_source_type> name: Database CDC Source # 请配置对应数据库 Source 的连接参数。 sink: type: kafka name: Kafka Sink properties.bootstrap.servers: <kafka_bootstrap_servers> topic: <target_topic> value.format: debezium-json sink.add-tableId-to-header-enabled: true pipeline: name: Database CDC to Kafka parallelism: 2
预期结果。 Pipeline 将数据库 INSERT、UPDATE、DELETE 序列化为 Debezium JSON,并写入指定 Topic。
source: type: kafka name: Kafka Source properties.bootstrap.servers: <source_kafka_bootstrap_servers> properties.group.id: <consumer_group_id> topic-pattern: <source_topic_pattern> value.format: debezium-json table-id.dynamic.from: value table-id.dynamic.fields: source.db,source.table scan.startup.mode: group-offsets sink: type: kafka name: Kafka Sink properties.bootstrap.servers: <target_kafka_bootstrap_servers> value.format: debezium-json pipeline: name: Kafka CDC Relay Pipeline parallelism: 2
预期结果。 Source 解析上游 CDC 消息,Sink 按 TableId 将消息写入下游 Kafka Topic。
Source 的首次启动位置由 scan.startup.mode 决定。启用 Checkpoint 后,Checkpoint 会保存 Kafka Offset;作业恢复时从已保存的位置继续读取,而不是重新应用首次启动位置,从而提供 Exactly-Once Source 处理能力。端到端一致性仍取决于下游 Sink;Kafka Sink 的交付语义为 At-Least-Once,消费端应基于业务主键或事件唯一标识实现幂等处理。
Kafka 只保证单 Partition 内有序。Source 并行读取多个 Partition;Sink 的 Key 和分区策略决定同一业务 Key 是否进入同一 Partition。调整 Pipeline 并行度或 Topic 分区数前,应评估顺序和数据倾斜。
使用 hash-by-key 时,Kafka Key 由 TableId 和主键字段组成。无主键表的 Key 仅包含 TableId,因此同一张表的记录会基于相同 Key 哈希到同一 Kafka Partition;该配置不会因缺少主键而报错,也不会回退为 all-to-zero。
Source 按 TableId 独立维护 Schema。新增字段会扩展 Schema,缺失字段不会自动删除。Sink 的 Schema 输出取决于上游 Source 是否产生 Schema Change,以及 Sink 的 value.format 和 debezium-json.include-schema.enabled 配置;Sink 不会独立推断上游表结构。
现象 | 可能原因 | 处理建议 |
|---|---|---|
Source 没有消费消息 | Topic/Pattern 不匹配、Group Offset 已在末尾、网络或认证失败。 | 检查 Topic、Group、启动模式、Broker 连通性和 Consumer 权限。 |
Source 解析失败 |
| 抽样检查原始消息,并核对格式和动态字段 Path。 |
TableId 不符合预期 | 动态字段数量、顺序或值不正确。 | 验证 Path 返回字符串,并按 1/2/3 字段规则检查结果。 |
Sink 没有写入 | 目标 Topic 未提前创建、Topic 权限不足或 Producer 参数错误。 | 创建目标 Topic,并检查 Producer 权限和连接参数。 |
同一主键消息乱序 | 相同 Key 被写入不同 Partition,或上游本身跨 Partition。 | 配置稳定 Key 和分区策略,并检查 Topic 分区与并行度。 |
字段类型不断扩展 | 同一 TableId 的消息字段类型不稳定。 | 统一上游消息 Schema,或在写入 Kafka 前规范字段类型。 |
WAL/数据库事件已产生但 Kafka 无消息 | Source 与 Sink 之间反压、Checkpoint 或 Kafka Producer 异常。 | 依次检查 Pipeline 拓扑、Checkpoint、Producer 指标和异常日志。 |
properties.* 配置 PLAINTEXT、SASL_PLAINTEXT、SASL_SSL 和 SSL;具体参数须与 Kafka 集群认证方式一致。ignore-parse-errors 会跳过错误消息,可能造成数据缺失;生产环境开启前应建立死信或审计方案。latest-offset 会跳过已有消息,earliest-offset 可能读取大量历史数据。