You need to enable JavaScript to run this app.
文档中心
流式计算 Flink版

流式计算 Flink版

复制全文
下载 pdf
CDC Connector 参考
Kafka
复制全文
下载 pdf
Kafka

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 版

  • 1.16
  • 1.20

Source 和 Sink 均为平台内置 Connector,无需额外上传 JAR。

Flink CDC

3.4.0 及以上

创建作业时请选择 3.4.0 或更高版本。

Kafka

支持火山引擎 Kafka、自建 Kafka,以及兼容 Kafka 协议的服务。

特色功能

  • Source 与 Sink 双向支持:同一 Connector 可用于消费 CDC 消息或将 Pipeline 变更写入 Kafka。
  • 多格式 CDC 消息:Source 支持 JSON、Debezium JSON 和 Canal JSON。
  • 灵活的 TableId:Source 可以固定 TableId,也可以从 Key 或 Value 的嵌套字段动态提取。
  • 灵活的 Topic 路由:Sink 可以固定写入单一 Topic,或按上游 TableId 分 Topic 写入。
  • Schema 推导与序列化:Source 根据消息字段合并 Schema;Sink 将 Pipeline 数据类型转换为 CDC JSON。

前提条件

  1. 已开通火山引擎流式计算 Flink 版,并创建 Flink CDC Pipeline 作业。
  2. 已准备 Kafka 集群,并提前创建 Source 订阅的 Topic 和 Sink 写入的目标 Topic。
  3. Flink 资源池与 Kafka Broker 网络可达;跨 VPC、公网或专线场景已配置相应网络。
  4. Kafka 用户具有 Source 所需的 Topic 读取、Consumer Group 权限,或 Sink 所需的目标 Topic 写入权限。
  5. 已在 Flink 中创建 Kafka 用户名、密码、证书等加密变量,不在 YAML 中填写明文凭证。
  6. Source 场景已确认消息格式、TableId 获取方式和启动位点;Sink 场景已确认 Topic 路由、分区策略和输出格式。
  7. 需要故障恢复和一致性保证时,已为 Pipeline 配置可用的 Checkpoint。

使用限制

  • Source 的 topictopic-pattern 必须二选一,不能同时配置。
  • Source 必须配置一种 TableId 解析方式:table-id.fixed.name,或同时配置 table-id.dynamic.fromtable-id.dynamic.fields
  • 固定 TableId 与动态 TableId 配置不能同时使用。
  • table-id.dynamic.fields 支持 1~3 个逗号分隔的 JSON Path,目标字段值必须为字符串。
  • Source 使用 latest-offset 时不会补读作业启动前的消息;使用 timestampspecific-offsets 前应确认位点仍在 Kafka 保留范围内。
  • Sink 配置固定 topic 后,所有表的消息都会写入同一个 Topic,需要通过消息中的 TableId 信息区分来源。
  • Kafka 只能保证单 Partition 内的消息顺序。需要同一主键或同一表有序时,应合理配置 Key、分区策略、Topic 分区数和 Pipeline 并行度。
  • Sink 不会代替用户创建目标 Topic。固定 Topic、按 TableId 生成的 Topic,以及映射规则中的 Topic 均须提前创建。

语法结构

用作 Source

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

sink:
  type: kafka
  name: Kafka Sink
  properties.bootstrap.servers: <kafka_bootstrap_servers>
  value.format: debezium-json
  topic: <target_topic>

参数说明

Source 连接与订阅参数

参数

说明

数据类型

是否必填

默认值

备注

type

Pipeline Connector 类型。

STRING

固定为 kafka

name

Source 名称。

STRING

用于作业拓扑和日志识别。

properties.bootstrap.servers

Kafka Broker 接入地址。

STRING

多个地址使用英文逗号分隔。

properties.group.id

Kafka Consumer Group ID。

STRING

不同业务链路建议使用不同 Group ID。

topic

订阅的 Topic。

STRING / LIST

条件必填

topic-pattern 二选一。

topic-pattern

订阅 Topic 的正则表达式。

STRING

条件必填

topic 二选一。

scan.topic-partition-discovery.interval

新 Topic 或 Partition 的发现间隔。

DURATION

5 min

设置为 0 关闭动态发现。

properties.*

Kafka Consumer 原生参数。

STRING

去除 properties. 前缀后传给 Kafka Client。

Source 格式与 TableId 参数

参数

说明

数据类型

是否必填

默认值

备注

key.format

Kafka Key 解析格式。

STRING

none

支持 nonejson

key.fields-prefix

Key 字段前缀。

STRING

用于避免 Key 与 Value 字段重名。

value.format

Kafka Value 解析格式。

STRING

json

建议使用debezium-jsoncanal-json解析格式。

value.fields-prefix

Value 字段前缀。

STRING

用于避免 Key 与 Value 字段重名。

table-id.fixed.name

为所有消息设置固定 TableId。

STRING

条件必填

与动态 TableId 配置二选一。

table-id.dynamic.from

动态 TableId 的提取位置。

STRING

条件必填

支持 keyvalue

table-id.dynamic.fields

动态 TableId 的字段 Path。

STRING

条件必填

支持 1~3 个逗号分隔 Path。

Source 启动、元数据与限速参数

参数

说明

数据类型

是否必填

默认值

备注

scan.startup.mode

Kafka 消费起始位点。

STRING

group-offsets

支持 earliest-offsetlatest-offsetgroup-offsetstimestampspecific-offsets

scan.startup.timestamp-millis

时间戳启动位点。

LONG

条件必填

timestamp 模式必填,单位为毫秒。

scan.startup.specific-offsets

指定 Partition Offset。

STRING

条件必填

示例:partition:0,offset:42;partition:1,offset:300

metadata.list

透传到下游的 Kafka Metadata。

STRING

支持 topicpartitionoffsettimestampheadersleader-epochtimestamp-type

rate-limit-num

Source 限速值。

LONG

单位为记录数。
例如,该参数设置为 1000000,表示整个算子的总 OPS 不超过 100W。

Debezium JSON Source 参数

参数

说明

数据类型

是否必填

默认值

备注

debezium-json.ignore-parse-errors

是否跳过无法解析的字段和消息。

BOOLEAN

false

开启后错误消息可能被丢弃。

debezium-json.infer-schema.primitive-as-string

是否将原始类型统一推导为 STRING。

BOOLEAN

false

debezium-json.include-schema.enabled

输入消息是否包含 Kafka Connect Schema。

BOOLEAN

false

应与上游消息结构保持一致。

Canal JSON Source 参数

参数

说明

数据类型

是否必填

默认值

备注

canal-json.ignore-parse-errors

是否跳过无法解析的字段和消息。

BOOLEAN

false

开启后错误消息可能被丢弃。

canal-json.infer-schema.primitive-as-string

是否将原始类型统一推导为 STRING。

BOOLEAN

false

canal-json.infer-schema.strategy

Schema 推导策略。

STRING

AUTO

支持 AUTOMYSQL_TYPE

canal-json.mysql.treat-tinyint1-as-boolean.enabled

是否将 MySQL TINYINT(1) 映射为 BOOLEAN。

BOOLEAN

true

仅在 MYSQL_TYPE 策略下生效。

Sink 参数

参数

说明

数据类型

是否必填

默认值

备注

type

Pipeline Connector 类型。

STRING

固定为 kafka

name

Sink 名称。

STRING

用于作业拓扑和日志识别。

properties.bootstrap.servers

Kafka Broker 接入地址。

STRING

多个地址使用英文逗号分隔。

value.format

Kafka Value 序列化格式。

STRING

debezium-json

支持 debezium-jsoncanal-json

topic

固定写入的 Topic。

STRING

不配置时按 TableId 生成 Topic。

sink.add-tableId-to-header-enabled

是否将 TableId 写入 Header。

BOOLEAN

false

Header Key 为 namespaceschemaNametableName

sink.custom-header

自定义 Kafka Header。

STRING

示例:key1:value1,key2:value2

properties.*

Kafka Producer 原生参数。

STRING

去除 properties. 前缀后传给 Kafka Producer。

Sink 路由与格式参数

参数

说明

数据类型

是否必填

默认值

备注

partition.strategy

Kafka 分区策略。

STRING

all-to-zero

all-to-zero 将所有记录写入 0 号分区;hash-by-key 按主键哈希分发。无主键时 Kafka Key 仅包含 TableId,同一张表的记录会路由到同一分区。

key.format

Kafka Key 序列化格式。

STRING

json

支持 jsoncsv

sink.tableId-to-topic.mapping

TableId 到 Topic 的映射。

STRING

多个映射用分号分隔;映射到的 Topic 须提前创建。

debezium-json.include-schema.enabled

Debezium JSON 是否携带 Schema。

BOOLEAN

false

仅适用于 debezium-json

TableId 与 Topic 路由

Source 固定 TableId

单表 Topic 可以配置固定 TableId:

table-id.fixed.name: <schema_name>.<table_name>

Source 动态 TableId

动态 TableId 从 Kafka Key 或 Value 中提取:

table-id.dynamic.from: value
table-id.dynamic.fields: source.db,source.table

字段数量决定生成方式:

字段数量

生成结果

1

将字段值解析为完整 TableId。

2

生成 schema.table

3

生成 namespace.schema.table

Sink Topic 路由

  • 配置 topic:所有消息写入固定 Topic。
  • 未配置 topic:默认根据上游 TableId 生成 Topic;对应 Topic 须提前创建。
  • 可以使用 Pipeline Route 修改 TableId,从而调整默认 Topic。
  • 可以使用 sink.tableId-to-topic.mapping 直接配置 TableId 到 Topic 的映射;映射到的 Topic 须提前创建。

Schema 与类型处理

Source Schema 合并

  • 消息出现新字段时,将字段加入 Schema,并产生新增可空列事件。
  • 消息缺少已有字段时,保留该字段并填充 NULL,不产生删除列事件。
  • 同名字段类型相同但精度不同时,合并为更高精度类型。
  • 同名字段类型不同时,合并为能够容纳两种输入的共同父类型。
  • 时间类型合并后保留输入中的最大精度。
  • Decimal 合并时保留较大的 Scale 和整数位数;总精度超过 38 时优先保留整数位。

Sink 输出格式

debezium-json 输出 beforeafteropsourcecanal-json 输出 olddatatypedatabasetablepkNames

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

以下示例从 Kafka 消费 Debezium JSON,根据 source.dbsource.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 传递给下游。

将数据库变更写入 Kafka

以下示例使用占位 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。

Kafka CDC 消息中转

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。

运行机制

Offset 与 Checkpoint

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

Schema 演进

Source 按 TableId 独立维护 Schema。新增字段会扩展 Schema,缺失字段不会自动删除。Sink 的 Schema 输出取决于上游 Source 是否产生 Schema Change,以及 Sink 的 value.formatdebezium-json.include-schema.enabled 配置;Sink 不会独立推断上游表结构。

验证与排障

验证方法

  1. 验证 Flink 资源池能够访问所有 Kafka Broker。
  2. Source 作业检查 Consumer Group、Topic/Pattern、启动位点和 TableId 解析结果。
  3. Sink 作业检查目标 Topic、分区、Key、Header 和 CDC JSON 内容。
  4. 对测试数据执行 INSERT、UPDATE、DELETE,确认下游消息格式与 Row Change 一致。
  5. 完成一次 Checkpoint 后受控重启,确认 Offset 恢复和消息重复情况符合产品语义。
  6. 新增字段或改变字段精度,验证 Source Schema 合并和 Sink 输出结果。

常见问题

现象

可能原因

处理建议

Source 没有消费消息

Topic/Pattern 不匹配、Group Offset 已在末尾、网络或认证失败。

检查 Topic、Group、启动模式、Broker 连通性和 Consumer 权限。

Source 解析失败

value.format 与实际消息不一致,或消息缺少 TableId 字段。

抽样检查原始消息,并核对格式和动态字段 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 指标和异常日志。

注意事项

  • 不要在 YAML 中填写 Kafka 密码、JAAS 配置或证书密钥,应使用 Flink 加密变量。
  • Source 和 Sink 可通过 properties.* 配置 PLAINTEXTSASL_PLAINTEXTSASL_SSLSSL;具体参数须与 Kafka 集群认证方式一致。
  • ignore-parse-errors 会跳过错误消息,可能造成数据缺失;生产环境开启前应建立死信或审计方案。
  • latest-offset 会跳过已有消息,earliest-offset 可能读取大量历史数据。
  • Sink 使用的固定 Topic、TableId 默认路由 Topic 和映射 Topic 均须提前创建。
  • Kafka Sink 为 At-Least-Once,故障恢复可能产生重复消息,消费端应设计幂等处理。
  • 固定 Topic 汇聚多表时,必须保留可靠的 TableId 信息,否则消费者无法区分来源表。
  • 使用 Canal JSON 或 Debezium JSON 时,上下游格式配置必须一致。

相关文档

最近更新时间:2026.07.31 12:00:49
这个页面对您有帮助吗?
有用
有用
无用
无用