变形

解码逻辑解码消息内容

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

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

解码逻辑解码消息内容

DecodeLogicalDecodingMessageContent SMT 可将 PostgreSQL 逻辑解码消息的二进制内容转换为结构化形式。

当 Debezium PostgreSQL 连接器捕获逻辑解码消息时,它会向 Kafka 发出消息事件记录。默认情况下,这些消息记录中的 content 字段包含编码后的二进制数据。为了便于其他 Kafka 消费者处理 PostgreSQL 事件消息,你可以使用 DecodeLogicalDecodingMessageContent SMT 解码原始消息的二进制内容,并将其转换为更易于消费的格式。

你还可以将该 SMT 与其他 SMT 配合使用,例如 Debezium Outbox Event Router。

示例

要使 Debezium PostgreSQL 连接器能够解码消息事件中的二进制内容,请将 DecodeLogicalDecodingMessageContent SMT 添加到该连接器的 Kafka Connect 配置中,如下例所示:

"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
...
"transforms": "decodeLogicalDecodingMessageContent",
"transforms.decodeLogicalDecodingMessageContent.type": "io.debezium.connector.postgresql.transforms.DecodeLogicalDecodingMessageContent",
"key.converter": "org.apache.kafka.connect.json.JsonConverter",
"key.converter.schemas.enable": false,
"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"value.converter.schemas.enable": "false",
...

以下示例展示了事件记录在应用转换前后的 key 与 value。

示例 1. 应用 DecodeLogicalDecodingMessageContent SMT 的效果

处理前

SMT 处理记录之前的事件 key

{
    "prefix": "test-prefix"
}

SMT 处理记录之前的事件值

{
    "op": "m",
    "ts_ms": 1723115240065,
    "source": {
        "version": "3.0.0-SNAPSHOT",
        "connector": "postgresql",
        "name": "connector-name",
        "ts_ms": 1723115239782,
        "snapshot": "false",
        "db": "source-db",
        "sequence": "[\"26997744\",\"26997904\"]",
        "ts_us": 1723115239782690,
        "ts_ns": 1723115239782690000,
        "schema": "",
        "table": "",
        "txId": 756,
        "lsn": 26997904,
        "xmin": null
    },
    "message": {
        "prefix": "test-prefix",
        "content": "eyJpZCI6IDEsICJpdGVtIjogIkRlYmV6aXVtIGluIEFjdGlvbiIsICJzdGF0dXMiOiAiRU5URVJFRCIsICJxdWFudGl0eSI6IDIsICJ0b3RhbFByaWNlIjogMzkuOTh9"
    }
}

处理后

SMT 处理记录后的事件键

null

SMT 处理记录后的事件值

{
    "op": "c",
    "ts_ms": 1723115415729,
    "source": {
        "version": "3.0.0-SNAPSHOT",
        "connector": "postgresql",
        "name": "connector-name",
        "ts_ms": 1723115415640,
        "snapshot": "false",
        "db": "source-db",
        "sequence": "[\"26717416\",\"26717576\"]",
        "ts_us": 1723115415640161,
        "ts_ns": 1723115415640161000,
        "schema": "",
        "table": "",
        "txId": 745,
        "lsn": 26717576,
        "xmin": null
    },
    "after": {
        "id": 1,
        "item": "Debezium in Action",
        "status": "ENTERED",
        "quantity": 2,
        "totalPrice": 39.98
    }
}

在上面的示例中,该 SMT 对原始事件记录应用了以下变更:

  • 删除原始逻辑解码消息中包含 prefix 字段的键("prefix": "test-prefix")。
  • 将 op 字段的值从 m(message)转换为 c(create),从而将事件类型从 message 改为 INSERT。
  • 用 after 字段替换 message 字段,其中包含逻辑解码消息的解码内容。

在 SMT 应用这些变更之后,下游消费者或其他 SMT(例如 Debezium Outbox Event Router)就能更方便地处理该记录。

配置选项

下表列出了可与 DecodeLogicalDecodingMessageContent SMT 一起使用的配置选项。

属性类型默认值说明
fields.null.includebooleanfalse指定解码过程如何处理源消息中值为 null 的字段。默认情况下,该转换会删除值为 null 的字段。

表 1. DecodeLogicalDecodingMessageContent SMT 配置选项

评论

登录后参与评论

正在加载评论…