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:调整兼容性配置
- 移除
type.name: "_doc"参数:Elasticsearch 7.x已默认使用_doc作为类型,显式配置可能引发兼容性问题。 - 临时关闭压缩排查:若不确定压缩配置是否兼容,可先设置
connection.compression: "false",测试是否能正常运行。
验证步骤
- 更新连接器配置并重启
- 发送包含JSON数组的测试消息到Kafka主题
- 检查Elasticsearch是否成功写入数据,且无报错
内容的提问来源于stack exchange,提问作者Emrahall
相关产品推荐
相关产品推荐

