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
相关产品推荐
相关产品推荐

