Elasticsearch Source Connector运行报错及数据读取验证咨询
Elasticsearch Source Connector 运行错误排查与数据读取确认
问题背景
已成功启动Elasticsearch Source Connector与Apache Kafka,连接器状态显示正常,但加载配置文件时出现错误,核心问题为No mapping found for [@timestamp] in order to sort on,同时需要确认Elasticsearch中的数据是否已被Kafka读取。
连接器状态
{ "name": "elastic-source", "connector": { "state": "RUNNING", "worker_id": "127.0.1.1:8083" }, "tasks": [ { "id": 0, "state": "RUNNING", "worker_id": "127.0.1.1:8083" } ], "type": "source" }
连接器配置文件(elasticsearch-source.properties)
name=elastic-source connector.class=com.github.dariobalinzo.ElasticSourceConnector tasks.max=1 es.host=127.0.0.1 es.port=1750 index.prefix=products topic.prefix=es_ topic=elastic-events
错误日志详情
核心错误提示
{"error":{"root_cause":[{"type":"query_shard_exception","reason":"No mapping found for [@timestamp] in order to sort on","index_uuid":"b575v16yTXmq5o2sk77zbA","index":"products"}],"type":"search_phase_execution_exception","reason":"all shards failed","phase":"can_match","grouped":true,"failed_shards":[{"shard":0,"index":"products","node":"asmTRFlgThS7kBU6yyCJzA","reason":{"type":"query_shard_exception","reason":"No mapping found for [@timestamp] in order to sort on","index_uuid":"b575v16yTXmq5o2sk77zbA","index":"products"}}]},"status":400}
完整堆栈日志
[2022-09-20 17:45:38,127] INFO [elastic-source|task-0] fetching from products (com.github.dariobalinzo.task.ElasticSourceTask:201) [2022-09-20 17:45:38,128] INFO [elastic-source|task-0] found last value Cursor{primaryCursor='null', secondaryCursor='null'} (com.github.dariobalinzo.task.ElasticSourceTask:203) [2022-09-20 17:45:38,129] WARN [elastic-source|task-0] request [POST http://localhost:1750/products/_search?typed_keys=true&max_concurrent_shard_requests=5&search_type=query_then_fetch&batched_reduce_size=512] returned 1 warnings: [299 Elasticsearch-7.15.0-79d65f6e357953a5b3cbcc5e2c7c21073d89aa29 "Elasticsearch built-in security features are not enabled. Without authentication, your cluster could be accessible to anyone."] (org.elasticsearch.client.RestClient:72) [2022-09-20 17:45:38,129] ERROR [elastic-source|task-0] error (com.github.dariobalinzo.task.ElasticSourceTask:217) ElasticsearchStatusException[Elasticsearch exception [type=search_phase_execution_exception, reason=all shards failed]] at org.elasticsearch.rest.BytesRestResponse.errorFromXContent(BytesRestResponse.java:178) at org.elasticsearch.client.RestHighLevelClient.parseEntity(RestHighLevelClient.java:2484) at org.elasticsearch.client.RestHighLevelClient.parseResponseException(RestHighLevelClient.java:2461) at org.elasticsearch.client.RestHighLevelClient.internalPerformRequest(RestHighLevelClient.java:2184) at org.elasticsearch.client.RestHighLevelClient.performRequest(RestHighLevelClient.java:2137) at org.elasticsearch.client.RestHighLevelClient.performRequestAndParseEntity(RestHighLevelClient.java:2105) at org.elasticsearch.client.RestHighLevelClient.search(RestHighLevelClient.java:1367) at com.github.dariobalinzo.elastic.ElasticRepository.executeSearch(ElasticRepository.java:176) at com.github.dariobalinzo.elastic.ElasticRepository.searchAfter(ElasticRepository.java:90) at com.github.dariobalinzo.task.ElasticSourceTask.poll(ElasticSourceTask.java:205) at org.apache.kafka.connect.runtime.WorkerSourceTask.poll(WorkerSourceTask.java:305) at org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:249) at org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:188) at org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:243) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Suppressed: org.elasticsearch.client.ResponseException: method [POST], host [http://localhost:1750], URI [/products/_search?typed_keys=true&max_concurrent_shard_requests=5&search_type=query_then_fetch&batched_reduce_size=512], status line [HTTP/1.1 400 Bad Request]
当前Kafka主题列表
__consumer_offsets connect-configs connect-offset connect-offsets connect-status elastic-events
解决方案与数据确认方法
1. 解决@timestamp映射缺失问题
该Elasticsearch Source Connector默认使用@timestamp字段作为排序和增量同步的依据,而你的products索引没有这个字段的映射,因此报错。可以通过以下两种方式解决:
- 方式一:给products索引添加@timestamp映射
执行Elasticsearch请求添加字段映射:
如果数据中已有其他时间字段,也可以将该字段重命名或复制为PUT /products/_mapping { "properties": { "@timestamp": { "type": "date" } } }@timestamp。 - 方式二:修改连接器配置,指定自定义时间字段
在elasticsearch-source.properties中添加配置项,指定你索引中已存在的时间字段作为同步依据,例如:
修改后重启连接器即可。timestamp.field=create_time # 替换为你实际的时间字段名
2. 确认Elasticsearch数据是否被Kafka读取
可以通过以下几种方式验证:
- 查看消费者组偏移量
执行Kafka命令查看连接器对应的消费者组偏移量(默认组名为connect-elastic-source):
如果kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group connect-elastic-sourceCURRENT-OFFSET大于0,说明已有数据被同步到Kafka。 - 直接消费主题消息
使用控制台消费者查看elastic-events主题的内容:
若能看到Elasticsearch中的数据,说明同步成功。kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic elastic-events --from-beginning - 检查连接器偏移量存储
查看Kafka的connect-offsets主题,搜索与elastic-source相关的条目,确认是否有记录同步的偏移位置。
内容的提问来源于stack exchange,提问作者Sarindra Thérèse
相关产品推荐
相关产品推荐

