Elasticsearch未从Logstash获取全部文档问题排查求助
文档导入不稳定问题排查建议
环境信息
- Logstash:8.6.0
- Elasticsearch:7.17.8
测试场景
预期向Elasticsearch导入28000个文档,重复执行以下流程10次:
- 删除Elasticsearch目标索引
- 启动Logstash
- 等待Pipeline终止,确认无报错/警告日志
- 在Kibana中核查文档数量
- 检查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?vAPI获取实时文档数,排除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
相关产品推荐
相关产品推荐

