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

Debezium Outbox Router序列化异常求助:PostgreSQL转换器适配问题

Debezium Outbox Router处理PostgreSQL JSON payload问题解决方法

问题场景

在PostgreSQL中使用Debezium Outbox Router时,payload列存储如下JSON数据:

{"userId": 107385,"chatId": "beb8faec-b75f-4eca-ace0-57b8621c7ca0","fromEcho": true,"results": [{"key": "agentSwitch","value": "Evet","type": "SWITCH"},{"key": "agentStar","value": {"chips": [],"rate": 5},"type": "STARS"},{"key": "agentDescription","value": null,"type": "TEXT_AREA"}]}

使用org.apache.kafka.connect.json.JsonConverter时的错误

当配置value.converter为org.apache.kafka.connect.json.JsonConverter时,连接器抛出如下错误:

{"id": 0,"state": "FAILED","worker_id": "10.233.66.78:8083","trace": "org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler\n\tat org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:230)\n\tat org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:156)\n\tat org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.convertTransformedRecord(AbstractWorkerSourceTask.java:494)\n\tat org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.sendRecords(AbstractWorkerSourceTask.java:402)\n\tat org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.execute(AbstractWorkerSourceTask.java:367)\n\tat org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:204)\n\tat org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:259)\n\tat org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.run(AbstractWorkerSourceTask.java:77)\n\tat org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:237)\n\tat java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)\n\tat java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)\n\tat java.base/java.lang.Thread.run(Thread.java:829)\nCaused by: org.apache.kafka.connect.errors.DataException: Invalid type for STRUCT: class java.lang.String\n\tat org.apache.kafka.connect.json.JsonConverter.convertToJson(JsonConverter.java:680)\n\tat org.apache.kafka.connect.json.JsonConverter.convertToJsonWithoutEnvelope(JsonConverter.java:563)\n\tat org.apache.kafka.connect.json.JsonConverter.fromConnectData(JsonConverter.java:313)\n\tat org.apache.kafka.connect.storage.Converter.fromConnectData(Converter.java:67)\n\tat org.apache.kafka.connect.runtime.AbstractWorkerSourceTask.lambda$convertTransformedRecord$6(AbstractWorkerSourceTask.java:494)\n\tat org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:180)\n\tat org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:214)\n\t... 13 more"}

使用org.apache.kafka.connect.storage.StringConverter时的异常

切换为org.apache.kafka.connect.storage.StringConverter后,虽然能正常发送初始JSON,但当payload数组删除一项后(如下JSON):

{"userId": 107385,"chatId": "beb8faec-b75f-4eca-ace0-57b8621c7ca0","chatUserSegment": "SellerMeal","agentNickName": "","rating": 5,"comment": "","duration": 0,"orderId": 0,"timestamp": 1716448682,"platform": "IosV2","fromEcho": true,"results": [{"key": "agentStar","value": {"chips": [],"rate": 5},"type": "STARS"},{"key": "agentDescription","value": null,"type": "TEXT_AREA"}]}

Kafka中生成的内容变为非JSON格式的Struct字符串:

Struct{chatId=beb8faec-b75f-4eca-ace0-57b8621c7ca0,rating=5,userId=107385,comment=,orderId=0,results=[Struct{key=agentStar,type=STARS,value=Struct{rate=5}}, Struct{key=agentDescription,type=TEXT_AREA}],duration=0,fromEcho=true,platform=IosV2,timestamp=1716448682,agentNickName=,chatUserSegment=SellerMeal}

原连接器配置

{"connector.class": "io.debezium.connector.postgresql.PostgresConnector","transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter","slot.name": "slotname","tasks.max": "1","publication.name": "publication_name","database.history.kafka.topic": "topic.history","transforms": "outbox","slot.max.retries": "3","slot.retry.delay.ms": "5000","transforms.outbox.table.expand.json.payload": "true","topic.prefix": "chat-assistant","transforms.outbox.route.topic.replacement": "${routedByValue}","value.converter": "org.apache.kafka.connect.storage.StringConverter","key.converter": "org.apache.kafka.connect.storage.StringConverter","database.user": "user","database.dbname": "db","transforms.outbox.table.fields.additional.placement": "header:header:customFields","transforms.outbox.table.field.event.key": "event_id","transforms.outbox.table.json.payload.null.behavior": "optional_bytes","database.server.name": "name","transforms.outbox.route.by.field": "topic","plugin.name": "pgoutput","database.port": "5432","key.converter.schemas.enable": "false","database.hostname": "ip","database.password": "password","name": "connector_name","value.converter.schemas.enable": "false","table.include.list": "public.outbox","transforms.outbox.table.field.event.payload.id": "event_id"}

问题原因

核心问题在于transforms.outbox.table.expand.json.payload参数设置为true时,Outbox Router会将payload列的JSON数据解析为Kafka Connect的Struct对象:

  • 使用JsonConverter时,转换器无法正确处理Struct与字符串的类型转换,抛出Invalid type for STRUCT错误;
  • 使用StringConverter时,当payload结构发生变化(如数组元素减少),Struct对象会被直接序列化为其toString()形式,导致输出非JSON格式的字符串。

解决方案

1. 修改Outbox转换配置

将transforms.outbox.table.expand.json.payload设置为false,这样Outbox Router不会解析JSON payload,而是直接传递原始的JSON字符串。

2. 换回JsonConverter并保持schema禁用

将value.converter改回org.apache.kafka.connect.json.JsonConverter,同时保持value.converter.schemas.enable为false,确保转换器直接处理原始JSON字符串,无需schema信息。

修改后的连接器配置

{"connector.class": "io.debezium.connector.postgresql.PostgresConnector","transforms.outbox.type": "io.debezium.transforms.outbox.EventRouter","slot.name": "slotname","tasks.max": "1","publication.name": "publication_name","database.history.kafka.topic": "topic.history","transforms": "outbox","slot.max.retries": "3","slot.retry.delay.ms": "5000","transforms.outbox.table.expand.json.payload": "false","topic.prefix": "chat-assistant","transforms.outbox.route.topic.replacement": "${routedByValue}","value.converter": "org.apache.kafka.connect.json.JsonConverter","key.converter": "org.apache.kafka.connect.storage.StringConverter","database.user": "user","database.dbname": "db","transforms.outbox.table.fields.additional.placement": "header:header:customFields","transforms.outbox.table.field.event.key": "event_id","transforms.outbox.table.json.payload.null.behavior": "optional_bytes","database.server.name": "name","transforms.outbox.route.by.field": "topic","plugin.name": "pgoutput","database.port": "5432","key.converter.schemas.enable": "false","database.hostname": "ip","database.password": "password","name": "connector_name","value.converter.schemas.enable": "false","table.include.list": "public.outbox","transforms.outbox.table.field.event.payload.id": "event_id"}

验证效果

修改配置后,无论payload的JSON结构如何变化,Kafka中都会收到原始的JSON格式数据,同时不会出现类型转换错误。

内容的提问来源于stack exchange,提问作者eren arslan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 14:57:02