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

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可成功导入

排查建议与解决方案

  1. 修复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
    
  2. 修正ES安全协议配置
    当前连接URL使用https但安全协议设为PLAINTEXT,存在不匹配问题,修改为:

    elastic.security.protocol=SSL
    
  3. 规范第三方数据格式
    从测试结果看,多行JSON无法批量导入ES,需确保第三方生产者写入的每条Kafka消息是单行合法JSON。若无法修改生产者,可添加Kafka Connect转换插件处理多行数据。

  4. 优化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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 05:20:54