变形

活动变更

qianmoQqianmoQ· 更新于 2026-09-28· 阅读 9 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

事件记录变更

此单一消息转换(SMT)仅适用于 SQL 数据库连接器。

Debezium 的数据变更事件结构复杂,提供了大量信息。然而,在某些情况下,下游消费者在处理 Debezium 变更事件消息之前,还需要了解由原始数据库变更所引起的字段级变化的额外信息。为了在事件消息中补充数据库操作如何修改源数据库中字段的详细信息,Debezium 提供了 ExtractChangedRecordState 单一消息转换(SMT)。

事件变更转换是一种 Kafka Connect SMT。

变更事件结构

Debezium 生成的数据变更事件结构复杂。每个事件由以下部分组成:

  • 元数据,包括但不限于以下类型:

    • 更改数据的操作类型。
    • 源信息,例如发生变更的数据库和表的名称。
    • 标识变更发生时间的时间戳。
    • 可选的事务信息。
  • 变更前的行数据。

  • 变更后的行数据。

以下示例展示了一个典型的 Debezium UPDATE 变更事件结构的部分内容:

{
    "op": "u",
    "source": {
        ...
    },
    "ts_ms" : "...",
    "ts_us" : "...",
    "ts_ns" : "...",
    "before" : {
        "field1" : "oldvalue1",
        "field2" : "oldvalue2"
    },
    "after" : {
        "field1" : "newvalue1",
        "field2" : "newvalue2"
    }
}

更多关于变更事件结构的详细信息,请参阅各连接器的文档。

前面示例中消息的复杂格式提供了关于源数据库中所发生变更的详细信息。但是,这种格式可能并不适合某些下游消费者。Sink 连接器或 Kafka 生态系统的其他部分可能期望消息能够明确标识出数据库操作修改了哪些字段、又保持了哪些字段不变。ExtractChangedRecordState 单消息转换(SMT)会向变更事件消息添加标头,用于标识数据库操作修改的字段以及保持不变的字段。

行为

事件变更 SMT 会从 Kafka 记录中的 Debezium UPDATE 变更事件里提取 before 和 after 字段。该转换会检查 before 和 after 事件状态结构,以识别被操作修改的字段以及保持不变的字段。根据连接器的配置,该转换随后会生成一个修改后的事件消息,添加消息标头以列出已变更的字段、未变更的字段,或两者都列出。如果事件表示的是 INSERT、DELETE 或 READ,该单消息转换会将配置的标头添加为空列表,因为不存在已变更或未变更的字段。

你可以为 Debezium 连接器配置事件变更 SMT,也可以为消费 Debezium 连接器所发出消息的 sink 连接器进行配置。如果你希望 Apache Kafka 保留完整的原始 Debezium 变更事件,则应为 sink 连接器配置事件变更 SMT。将该 SMT 应用于源连接器还是 sink 连接器,取决于你的具体用例。

根据你的用例,你可以配置该转换来修改原始消息,执行以下一项或两项任务:

  • 通过在用户配置的 header.changed.name 标头中列出相应字段,标识 UPDATE 事件所变更的字段。
  • 通过在用户配置的 header.unchanged.name 标头中列出相应字段,标识 UPDATE 事件未变更的字段。

配置

你可以通过在连接器的配置中添加该 SMT 配置详情,来为 Kafka Connect 源连接器或 sink 连接器配置 Debezium 事件变更 SMT。若要获得不添加任何标头的默认行为,只需将该转换添加到连接器配置中即可,如下例所示:

transforms=changes,...
transforms.changes.type=io.debezium.transforms.ExtractChangedRecordState

与任何 Kafka Connect 连接器配置一样,你可以将 transforms= 设置为多个以逗号分隔的 SMT 别名,并按你希望 Kafka Connect 应用这些 SMT 的顺序排列。

下面示例中的连接器配置为事件变更 SMT 设置了若干选项:

transforms=changes,...
transforms.changes.type=io.debezium.transforms.ExtractChangedRecordState
transforms.changes.header.changed.name=Changed
transforms.changes.header.unchanged.name=Unchanged

header.changed.name

用于存储由数据库操作所更改字段的逗号分隔列表的 Kafka 消息标头名称。

header.unchanged.name

用于存储数据库操作后仍保持不变字段的逗号分隔列表的 Kafka 消息标头名称。

自定义配置

连接器可能会发出多种类型的事件消息(心跳消息、墓碑消息,或有关事务和模式变更的元数据消息)。要将转换仅应用于部分事件,你可以定义一条 SMT 谓词语句,以便有选择地对特定事件应用转换。

有选择地应用事件变更转换的选项

除了 Debezium 连接器在数据库发生变更时发出的变更事件消息外,连接器还会发出其他类型的消息,包括心跳消息,以及有关模式变更和事务的元数据消息。由于这些其他消息的结构与该 SMT 设计用于处理的变更事件消息的结构不同,最好配置连接器以有选择地应用该 SMT,使其仅处理预期的数据变更消息。

有关如何有选择地应用 SMT 的更多信息,请参阅为转换配置 SMT 谓词。

配置选项

下表描述了可用于配置事件变更 SMT 的选项。

选项默认值描述
header.changed.name用于存储由数据库操作所更改字段的逗号分隔列表的 Kafka 消息标头名称。
header.unchanged.name用于存储数据库操作后仍保持不变字段的逗号分隔列表的 Kafka 消息标头名称。

表 1. 事件变更 SMT 配置选项说明

评论

登录后参与评论

正在加载评论…