集成

变更事件 SerDes

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

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

Debezium 事件反序列化

此功能目前处于孵化阶段,也就是说,确切的语义、配置选项等可能会根据我们收到的反馈在后续版本中发生变化。如果您在使用这些 SerDes 时遇到任何问题,请告诉我们。

Debezium 以复杂消息结构的形式生成数据变更事件。该消息随后会由所配置的 Kafka Connect 转换器进行序列化,而将其反序列化为逻辑消息则是消费者的职责。为此,Kafka 使用了所谓的 SerDes。

Debezium 提供了 SerDes(io.debezium.serde.DebeziumSerdes),以简化消费者的反序列化过程,无论该消费者是 Kafka Streams 管道还是普通的 Kafka 消费者。

JSON SerDe

JSON SerDe 将以 JSON 编码的变更事件反序列化并转换为 Java 类。在内部,这是通过 Jackson Databind 实现的。

消费者使用

final Serde<MyType> serde = DebeziumSerdes.payloadJson(MyType.class);

消费者随后会收到逻辑 Java 类型 MyType,其字段由该 JSON 消息初始化。这一规则同时适用于键和值。也可以使用 Integer 之类的普通 Java 类型,例如当键只由单个 INT 字段组成时。

当 Kafka Connect 使用 JSON 转换器时,它通常提供两种运行模式——带 schema 或不带 schema。如果使用 schema,则消息的形式如下:

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

而不带 schema 时,结构看起来更像这样:

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

反序列化器的行为由 from.field 配置选项驱动,并遵循以下规则:

  • 如果消息包含 schema,则只使用 payload

  • 如果反序列化键,则将键字段映射到目标类中

  • 如果反序列化值且其包含 Debezium 事件封装(envelope),那么:

    • 如果未设置 from.field,则将整个封装反序列化为目标类型
    • 否则,仅将所配置字段的内容反序列化并映射为目标类型,从而有效地将消息扁平化
  • 如果反序列化值且其已包含扁平化消息(即使用 事件扁平化 的 SMT 时),则将该扁平化记录映射为目标逻辑类型

配置选项

属性默认值说明
from.fieldN/A如果需要反序列化包含完整封装的消息,则留空;如果只需要变更之前或之后的数据值,则设置为 before/after。
unknown.properties.ignoredfalse决定遇到未知属性时是应静默忽略,还是应抛出运行时异常。

评论

登录后参与评论

正在加载评论…