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

Elasticsearch未从Logstash获取全部文档问题排查求助

文档导入不稳定问题排查建议

环境信息

  • Logstash:8.6.0
  • Elasticsearch:7.17.8

测试场景

预期向Elasticsearch导入28000个文档,重复执行以下流程10次:

  1. 删除Elasticsearch目标索引
  2. 启动Logstash
  3. 等待Pipeline终止,确认无报错/警告日志
  4. 在Kibana中核查文档数量
  5. 检查Elasticsearch日志,确认无报错/警告

执行结果

结果不稳定,10次执行中:

  • 2次成功导入全部28000个文档
  • 4次仅导入4002个文档(等待4小时后数量无变化)
  • 4次导入27713个文档

当前已排查情况

  • 已开启dead_letter_queue.enable,但DLQ为空;根据重试策略,400/404/409类错误会触发DLQ或警告,推测无此类错误
  • Logstash日志显示已完成所有查询,但无法确认是Logstash未发送全部文档,还是Elasticsearch丢弃了部分文档

附:相关日志与配置

Logstash日志片段

[2023-03-31T10:21:55,645][INFO ][logstash.inputs.jdbc     ][index][d9299291c383ba4e70e5116a1b01ba758087b0cff45119233811238f41d0be0d] (0.006227s) SELECT * FROM (
   SELECT
       *,
       'PERSON' + CAST(PERSON.Id AS VARCHAR(20)) AS uniqueid
       FROM PERSON) AS [T1] ORDER BY 1 OFFSET 14000 ROWS FETCH NEXT 1000 ROWS ONLY
[2023-03-31T10:21:56,427][INFO ][logstash.javapipeline    ][index] Pipeline terminated {"pipeline.id"=>"index"}

Logstash配置

input {
   jdbc {
       jdbc_connection_string => "#connectionstring"
       jdbc_user => "#username"
       jdbc_password => "#password"
       jdbc_driver_library => "./logstash-core/lib/jars/mssql-jdbc-11.2.3.jre18.jar"
       jdbc_driver_class => "com.microsoft.sqlserver.jdbc.SQLServerDriver"
       clean_run => true
       record_last_run => false
       jdbc_paging_enabled => true
       jdbc_page_size => 1000
       jdbc_fetch_size => 1000
       statement => "
       SELECT
       *,
       'PERSON' + CAST(PERSON.Id AS VARCHAR(20)) AS uniqueid
       FROM PERSON"
   }
   jdbc {
       jdbc_connection_string => "#connectionstring"
       jdbc_user => "#username"
       jdbc_password => "#password"
       jdbc_driver_library => "./logstash-core/lib/jars/mssql-jdbc-11.2.3.jre18.jar"
       jdbc_driver_class => "com.microsoft.sqlserver.jdbc.SQLServerDriver"
       clean_run => true
       record_last_run => false
       jdbc_paging_enabled => true
       jdbc_page_size => 1000
       jdbc_fetch_size => 1000
       statement => "
       SELECT
       *,
       'INVOICE' + CAST(INVOICE.Id AS VARCHAR(20)) AS uniqueid
       FROM INVOICE"
   }
}
output {
   elasticsearch {
       hosts => ["#elasticsearchurl"]
       index => "#elasticindex"
       action => "index"
       document_id => "%{uniqueid}"
   }
}

排查方向

1. 确认Logstash实际输出的文档数

  • 临时在Logstash的output中添加stdout { codec => rubydebug },将文档输出到控制台统计总数(数据量大时可测试小批量)
  • 启用metrics插件采集pipeline处理的文档总数:
    filter {
      metrics {
        meter => "documents"
        add_tag => "metrics"
      }
    }
    output {
      if "metrics" in [tags] {
        stdout {
          codec => line { format => "Total documents processed: %{[metrics][documents][count]}" }
        }
      }
    }
    
  • 将Logstash日志级别调整为debug,重点查看logstash.outputs.elasticsearch的日志,确认每个批量请求的发送数量和Elasticsearch的响应结果

2. 核查Elasticsearch的文档处理细节

  • 调用GET /_cat/count/#elasticindex?v API获取实时文档数,排除Kibana缓存影响
  • 检查Elasticsearch慢日志(配置文件中index.search.slowlog和index.indexing.slowlog),查看是否有批量写入超时或异常记录
  • 调用GET /_cat/nodes?v查看磁盘使用率disk.used_percent,若超过85%,Elasticsearch会触发只读模式阻止写入
  • 调用GET /_cat/thread_pool/bulk?v查看bulk线程池状态,若rejected数量大于0,说明bulk请求被拒绝

3. 排查JDBC输入的稳定性

  • 确认源数据库中PERSON和INVOICE表的总记录数是否固定,排除源数据变化导致的导入数量不一致
  • 启用JDBC输入的debug日志,查看是否存在分页查询中断、未执行完所有分页的情况
  • 检查Logstash的JVM内存配置(jvm.options中的-Xms和-Xmx),内存不足可能导致输入线程意外终止

4. 优化批量写入配置与网络排查

  • 调整elasticsearch输出的批量参数:
    elasticsearch {
      # 其他配置不变
      flush_size => 5000
      idle_flush_time => 30
      timeout => 60
      retry_on_conflict => 3
    }
    
  • 测试Logstash与Elasticsearch之间的网络连通性,排查是否存在丢包、超时情况
  • 确认Elasticsearch的http.max_content_length配置,是否限制了过大的bulk请求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 12:47:51