解码逻辑解码消息内容
解码逻辑解码消息内容
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 处理记录后的事件键
nullSMT 处理记录后的事件值
{
"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.include | boolean | false | 指定解码过程如何处理源消息中值为 null 的字段。默认情况下,该转换会删除值为 null 的字段。 |
表 1. DecodeLogicalDecodingMessageContent SMT 配置选项
评论
登录后参与评论
KnowForge