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

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-source
    
    如果CURRENT-OFFSET大于0,说明已有数据被同步到Kafka。
  • 直接消费主题消息
    使用控制台消费者查看elastic-events主题的内容:
    kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic elastic-events --from-beginning
    
    若能看到Elasticsearch中的数据,说明同步成功。
  • 检查连接器偏移量存储
    查看Kafka的connect-offsets主题,搜索与elastic-source相关的条目,确认是否有记录同步的偏移位置。

内容的提问来源于stack exchange,提问作者Sarindra Thérèse

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 00:05:23