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

Kafka Elasticsearch Sink连接器处理JSON数组报错求助

问题描述

使用Elasticsearch Sink连接器将Kafka主题消息同步至Elasticsearch时,主题包含JSON数组的情况下出现报错,不含数组的连接器运行正常。

版本信息

  • Elasticsearch: 7.17.13
  • Kafka: 2.6.0
  • Elasticsearch连接器库:Confluent Kafka Connect Elasticsearch

连接器配置

curl --location --request POST 'localhost:8084/connectors' --header 'Content-Type: application/json' --data-raw '
{
  "name": "TEST",
  "config": {
    "connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
    "tasks.max": "1",
    "topics": "TEST",
    "key.ignore": "true",
    "schema.ignore": "true",
    "connection.compression": "true",
    "key.converter.schemas.enable": "false",
    "value.converter.schemas.enable": "false",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "connection.url": "http://localhost:9200",
    "connection.username" :"elastic",
    "connection.password":"password",
    "type.name": "_doc",
    "name": "TEST"
  }
}'

异常信息

org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.
    at org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:588)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:323)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:226)
    at org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:198)
    at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:185)
    at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:235)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.run(FutureTask.java:266)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)
Caused by: org.apache.kafka.connect.errors.ConnectException: Bulk request failed
    at io.confluent.connect.elasticsearch.ElasticsearchClient$1.afterBulk(ElasticsearchClient.java:443)
    at org.elasticsearch.action.bulk.BulkRequestHandler$1.onFailure(BulkRequestHandler.java:64)
    at org.elasticsearch.action.ActionListener$Delegating.onFailure(ActionListener.java:66)
    at org.elasticsearch.action.ActionListener$RunAfterActionListener.onFailure(ActionListener.java:350)
    at org.elasticsearch.action.ActionListener$Delegating.onFailure(ActionListener.java:66)
    at org.elasticsearch.action.bulk.Retry$RetryHandler.onFailure(Retry.java:123)
    at io.confluent.connect.elasticsearch.ElasticsearchClient.lambda$null$1(ElasticsearchClient.java:216)
    ... 5 more
    Caused by: org.apache.kafka.connect.errors.ConnectException: Failed to execute bulk request due to 'org.elasticsearch.common.compress.NotXContentException: Compressor detection can only be called on some xcontent bytes or compressed xcontent bytes' after 6 attempt(s)
        at io.confluent.connect.elasticsearch.RetryUtil.callWithRetries(RetryUtil.java:165)
        at io.confluent.connect.elasticsearch.RetryUtil.callWithRetries(RetryUtil.java:119)
        at io.confluent.connect.elasticsearch.ElasticsearchClient.callWithRetries(ElasticsearchClient.java:490)
        at io.confluent.connect.elasticsearch.ElasticsearchClient.lambda$null$1(ElasticsearchClient.java:210)
        ... 5 more
    Caused by: org.elasticsearch.common.compress.NotXContentException: Compressor detection can only be called on some xcontent bytes or compressed xcontent bytes
        at org.elasticsearch.common.compress.CompressorFactory.compressor(CompressorFactory.java:42)
        at org.elasticsearch.common.xcontent.XContentHelper.createParser(XContentHelper.java:76)
        at org.elasticsearch.client.RequestConverters.bulk(RequestConverters.java:226)
        at org.elasticsearch.client.RestHighLevelClient.internalPerformRequest(RestHighLevelClient.java:2167)
        at org.elasticsearch.client.RestHighLevelClient.performRequest(RestHighLevelClient.java:2137)
        at org.elasticsearch.client.RestHighLevelClient.performRequestAndParseEntity(RestHighLevelClient.java:2105)
        at org.elasticsearch.client.RestHighLevelClient.bulk(RestHighLevelClient.java:620)
        at io.confluent.connect.elasticsearch.ElasticsearchClient.lambda$null$0(ElasticsearchClient.java:212)
        at io.confluent.connect.elasticsearch.RetryUtil.callWithRetries(RetryUtil.java:158)

解决方案

报错核心原因是Elasticsearch Sink连接器默认期望每条Kafka消息是单个JSON对象,而你的消息是JSON数组,导致连接器无法正确解析内容,触发NotXContentException。以下是可行的解决方法:

方案1:用Kafka Connect Transform拆分数组

通过添加Transform将JSON数组拆分为独立的JSON对象消息,推荐使用jq Transform处理:

修改连接器配置,替换原value.converter并添加Transform参数:

"value.converter": "org.apache.kafka.connect.json.JsonConverter",
"transforms": "splitArray",
"transforms.splitArray.type": "io.confluent.connect.jq.JqTransform$Value",
"transforms.splitArray.query": ".[]"

该jq表达式会将数组中的每个元素拆分为单独的消息,供Sink连接器正常处理。

方案2:在生产者端拆分数组

如果可以修改Kafka生产者代码,在发送消息前将JSON数组拆分为多个独立的JSON对象,分别发送到Kafka主题,Sink连接器即可直接处理,无需额外配置。

方案3:调整兼容性配置

  1. 移除type.name: "_doc"参数:Elasticsearch 7.x已默认使用_doc作为类型,显式配置可能引发兼容性问题。
  2. 临时关闭压缩排查:若不确定压缩配置是否兼容,可先设置connection.compression: "false",测试是否能正常运行。

验证步骤
  1. 更新连接器配置并重启
  2. 发送包含JSON数组的测试消息到Kafka主题
  3. 检查Elasticsearch是否成功写入数据,且无报错

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 02:20:57