You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Debezium EventRouter转换报错:无法找到payload字段问题咨询

问题:Debezium EventRouter转换报错“Unable to find payload field payload in event”

在Docker容器中基于Postgres、Kafka和Debezium Connect实现Outbox Pattern与CDC功能时,当outbox表新增记录,配置EventRouter转换后触发错误“Unable to find payload field payload in event”;移除transforms相关配置后错误消失,但无法自定义发送到Kafka的消息格式。

错误日志

2022-11-25 21:01:27 2022-11-25 20:01:27,304 INFO   Postgres|coordinator|streaming  First LSN 'LSN{0/184D1E8}' received   [io.debezium.connector.postgresql.connection.WalPositionLocator]
2022-11-25 21:01:27 2022-11-25 20:01:27,307 INFO   Postgres|coordinator|streaming  WAL resume position 'LSN{0/184D1E8}' discovered   [io.debezium.connector.postgresql.PostgresStreamingChangeEventSource]
2022-11-25 21:01:27 2022-11-25 20:01:27,310 INFO   Postgres|coordinator|streaming  Connection gracefully closed   [io.debezium.jdbc.JdbcConnection]
2022-11-25 21:01:27 2022-11-25 20:01:27,351 INFO   Postgres|coordinator|streaming  Requested thread factory for connector PostgresConnector, id = coordinator named = keep-alive   [io.debezium.util.Threads]
2022-11-25 21:01:27 2022-11-25 20:01:27,351 INFO   Postgres|coordinator|streaming  Creating thread debezium-postgresconnector-coordinator-keep-alive   [io.debezium.util.Threads]
2022-11-25 21:01:27 2022-11-25 20:01:27,351 INFO   Postgres|coordinator|streaming  Processing messages   [io.debezium.connector.postgresql.PostgresStreamingChangeEventSource]
2022-11-25 21:01:27 2022-11-25 20:01:27,352 INFO   Postgres|coordinator|streaming  Message with LSN 'LSN{0/184D1E8}' arrived, switching off the filtering   [io.debezium.connector.postgresql.connection.WalPositionLocator]
2022-11-25 21:02:03 2022-11-25 20:02:03,508 INFO   ||  2 records sent during previous 00:03:51.693, last recorded offset of {server=coordinator} partition is {transaction_id=null, lsn_proc=25679872, lsn_commit=25679448, lsn=25679872, txId=742, ts_usec=1669406522764811}   [io.debezium.connector.common.BaseSourceTask]
2022-11-25 21:02:03 2022-11-25 20:02:03,528 ERROR  ||  WorkerSourceTask{id=apimessage-outbox-connector-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted   [org.apache.kafka.connect.runtime.WorkerTask]
2022-11-25 21:02:03 org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:223)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:149)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.TransformationChain.apply(TransformationChain.java:50)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.sendRecords(AbstractWorkerSourceTask.java:386)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.execute(AbstractWorkerSourceTask.java:354)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:189)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:244)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.run(AbstractWorkerSourceTask.java:72)
2022-11-25 21:02:03     at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
2022-11-25 21:02:03     at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
2022-11-25 21:02:03     at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
2022-11-25 21:02:03     at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
2022-11-25 21:02:03     at java.base/java.lang.Thread.run(Thread.java:829)
2022-11-25 21:02:03 Caused by: org.apache.kafka.connect.errors.ConnectException: Unable to find payload field payload in event
2022-11-25 21:02:03     at io.debezium.transforms.outbox.EventRouterDelegate.apply(EventRouterDelegate.java:132)
2022-11-25 21:02:03     at io.debezium.transforms.outbox.EventRouter.apply(EventRouter.java:25)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.TransformationChain.lambda$apply$0(TransformationChain.java:50)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:173)
2022-11-25 21:02:03     at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:207)
2022-11-25 21:02:03     ... 12 more
2022-11-25 21:02:03 2022-11-25 20:02:03,533 INFO   ||  Stopping down connector   [io.debezium.connector.common.BaseSourceTask]
2022-11-25 21:02:03 2022-11-25 20:02:03,538 INFO   Postgres|coordinator|streaming  Connection gracefully closed   [io.debezium.jdbc.JdbcConnection]
2022-11-25 21:02:03 2022-11-25 20:02:03,538 INFO   Postgres|coordinator|streaming  Finished streaming   [io.debezium.pipeline.ChangeEventSourceCoordinator]
2022-11-25 21:02:03 2022-11-25 20:02:03,539 INFO   Postgres|coordinator|streaming  Connected metrics set to 'false'   [io.debezium.pipeline.ChangeEventSourceCoordinator]
2022-11-25 21:02:03 2022-11-25 20:02:03,543 INFO   ||  Connection gracefully closed   [io.debezium.jdbc.JdbcConnection]
2022-11-25 21:02:03 2022-11-25 20:02:03,544 INFO   ||  [Producer clientId=connector-producer-apimessage-outbox-connector-0] Closing the Kafka producer with timeoutMillis = 30000 ms.   [org.apache.kafka.clients.producer.KafkaProducer]
2022-11-25 21:02:03 2022-11-25 20:02:03,548 INFO   ||  Metrics scheduler closed   [org.apache.kafka.common.metrics.Metrics]
2022-11-25 21:02:03 2022-11-25 20:02:03,548 INFO   ||  Closing reporter org.apache.kafka.common.metrics.JmxReporter   [org.apache.kafka.common.metrics.Metrics]
2022-11-25 21:02:03 2022-11-25 20:02:03,548 INFO   ||  Metrics reporters closed   [org.apache.kafka.common.metrics.Metrics]
2022-11-25 21:02:03 2022-11-25 20:02:03,548 INFO   ||  App info kafka.producer for connector-producer-apimessage-outbox-connector-0 unregistered   [org.apache.kafka.common.utils.AppInfoParser]

连接器创建请求

curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '{
"name": "apimessage-outbox-connector",
"config": {
  "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
  "topic.prefix": "coordinator",
  "slot.name": "coordinator",
  "tasks.max": "1",
  "database.hostname": "postgres",
  "database.port": "5432",
  "database.user": "postgres",
  "database.password": "admin",
  "database.dbname": "coordinator",
  "database.server.name": "localhost",
  "tombstones.on.delete": "false",
  "table.whitelist": "coordinator.outbox_event",
  "transforms": "outbox",
  "transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter"}
}'

解决方案

1. 检查Outbox表结构

EventRouter默认要求表中存在payload字段(推荐使用Postgres的JSONB类型存储事件内容),同时通常需要aggregateid、type等字段。如果你的表字段名与默认不符,需要在连接器配置中指定映射关系。

2. 配置自定义字段映射

如果你的outbox表使用非默认字段名(比如event_payload代替payload,entity_id代替aggregateid),在连接器配置中添加以下参数:

"transforms.outbox.field.payload": "event_payload",
"transforms.outbox.field.aggregate.id": "entity_id",
"transforms.outbox.field.event.type": "event_type"

完整配置示例:

{
"name": "apimessage-outbox-connector",
"config": {
  "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
  "topic.prefix": "coordinator",
  "slot.name": "coordinator",
  "tasks.max": "1",
  "database.hostname": "postgres",
  "database.port": "5432",
  "database.user": "postgres",
  "database.password": "admin",
  "database.dbname": "coordinator",
  "database.server.name": "localhost",
  "tombstones.on.delete": "false",
  "table.whitelist": "coordinator.outbox_event",
  "transforms": "outbox",
  "transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter",
  "transforms.outbox.field.payload": "your_payload_field_name",
  "transforms.outbox.field.aggregate.id": "your_aggregate_id_field",
  "transforms.outbox.field.event.type": "your_event_type_field"
}
}

3. 验证原始CDC事件结构

移除EventRouter配置后,查看Kafka中的原始消息,确认after节点下是否包含payload字段。原始消息结构示例:

{
  "after": {
    "id": "1",
    "payload": "{\"key\": \"value\"}",
    "aggregateid": "123",
    "type": "UserCreated"
  }
}

如果after中没有对应字段,说明表结构不符合要求,需调整表结构或配置字段映射。

4. 确认Debezium版本兼容性

建议使用Debezium 2.x及以上稳定版本,旧版本可能存在Outbox模式的字段兼容性问题。

内容的提问来源于stack exchange,提问作者Jacek Grzegorczyk

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.10 21:22:02