Kafka Connect Elasticsearch Sink JsonConverter序列化异常求助
问题排查:Kafka Connect Elasticsearch Sink 序列化错误
环境与背景
- 搭建了包含3个broker、3个controller和1个worker的Apache Kafka集群
- 通过Confluent Elasticsearch Sink插件同步第三方数据到Elasticsearch
- 第三方数据无Schema,示例格式如下:
{ "location": "Battery Tester1", "dc": "16.20V", "ac": "12.01V", "curr": " 0.00A", "temperature": "32.00C", "status": [ "Currently on AC power" ]} { "location": "Battery Tester2", "dc": "16.10V", "ac": "11.01V", "curr": " 2.00A", "temperature": "34.00C", "status": [ "Currently on AC power" ]} { "location": "Battery Tester3", "status": [ "Currently on AC power" ]}
配置文件
connect-standalone.properties
bootstrap.servers=kafbrk01-4:9092,kafbrk01-5:9092,kafbrk01-6:9092 config.storage.topic: es-connect-kafwrk01-configs offset.storage.topic: es-connect-kafwrk01-offsets status.storage.topic: es-connect-kafwrk01-status config.storage.replication.factor: -1 offset.storage.replication.factor: -1 status.storage.replication.factor: -1 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 offset.storage.file.filename=/tmp/connect.offsets offset.flush.interval.ms=10000 plugin.path=/opt/kafka/plugins
Elasticsearch Sink 配置(quickstart-elasticsearch.properties)
name=elasticsearch-sink connector.class=io.confluent.connect.elasticsearch.ElasticsearchSinkConnector tasks.max=1 topics=Power,Router,Gateway key.ignore=true connection.url=https://<FQDN to es01>:9200,https://<FQDN to es02>:9200,https://<FQDN to es03>:9200 connection.username=es_sink_connector_user connection.password=FakePasswordBecause! type.name=kafka-connect elastic.security.protocol = PLAINTEXT schema.ignore=true
错误日志
[2024-04-24 22:06:15,672] ERROR [elasticsearch-sink|task-0] WorkerSinkTask{id=elasticsearch-sink-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:212) org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:230) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execute(RetryWithToleranceOperator.java:156) at org.apache.kafka.connect.runtime.WorkerSinkTask.convertAndTransformRecord(WorkerSinkTask.java:533) at org.apache.kafka.connect.runtime.WorkerSinkTask.convertMessages(WorkerSinkTask.java:513) at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:349) at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:250) at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:219) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:204) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:259) at org.apache.kafka.connect.runtime.isolation.Plugins.lambda$withClassLoader$1(Plugins.java:236) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:840) Caused by: org.apache.kafka.connect.errors.DataException: Converting byte[] to Kafka Connect data failed due to serialization error: at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:333) at org.apache.kafka.connect.storage.Converter.toConnectData(Converter.java:91) at org.apache.kafka.connect.runtime.WorkerSinkTask.lambda$convertAndTransformRecord$3(WorkerSinkTask.java:533) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndRetry(RetryWithToleranceOperator.java:180) at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:214) ... 14 more Caused by: org.apache.kafka.common.errors.SerializationException: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'key': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false') at [Source: (byte[])"key"; line: 1, column: 4] at org.apache.kafka.connect.json.JsonDeserializer.deserialize(JsonDeserializer.java:69) at org.apache.kafka.connect.json.JsonConverter.toConnectData(JsonConverter.java:331) ... 18 more Caused by: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'key': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false') at [Source: (byte[])"key"; line: 1, column: 4] at com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:2391) at com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:745) at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._reportInvalidToken(UTF8StreamJsonParser.java:3635) at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._handleUnexpectedValue(UTF8StreamJsonParser.java:2734) at com.fasterxml.jackson.core.json.UTF8StreamJsonParser._nextTokenNotInObject(UTF8StreamJsonParser.java:902) at com.fasterxml.jackson.core.json.UTF8StreamJsonParser.nextToken(UTF8StreamJsonParser.java:794) at com.fasterxml.jackson.databind.ObjectMapper._readTreeAndClose(ObjectMapper.java:4703) at com.fasterxml.jackson.databind.ObjectMapper.readTree(ObjectMapper.java:3090) at org.apache.kafka.connect.json.JsonDeserializer.deserialize(JsonDeserializer.java:67) ... 19 more
已尝试操作
- Python程序写入数据时Sink正常工作,第三方数据触发上述错误
- 将key/value converter改为
org.apache.kafka.connect.storage.StringConverter,生成空索引并报org.elasticsearch.common.compress.NotXContentException错误 - 导出主题数据到文件,添加
{"create":{"_index":"<topicname>"}}后curl批量导入失败,单行JSON可成功导入
排查建议与解决方案
修复Key解析错误
错误核心是Kafka记录的Key字段不是合法JSON,即使配置了key.ignore=true,JsonConverter仍会尝试解析Key。需单独为Key配置字符串转换器:
修改connect-standalone.properties:key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false修正ES安全协议配置
当前连接URL使用https但安全协议设为PLAINTEXT,存在不匹配问题,修改为:elastic.security.protocol=SSL规范第三方数据格式
从测试结果看,多行JSON无法批量导入ES,需确保第三方生产者写入的每条Kafka消息是单行合法JSON。若无法修改生产者,可添加Kafka Connect转换插件处理多行数据。优化Connect存储配置
原配置中replication.factor=-1无效,针对3 broker集群,设置为3以保证高可用:config.storage.replication.factor=3 offset.storage.replication.factor=3 status.storage.replication.factor=3
内容的提问来源于stack exchange,提问作者Scott C
相关产品推荐
相关产品推荐

