添加复合ID后Logstash向Elasticsearch导入数据未正常索引求助
问题诊断与解决方案
核心问题分析
出现docs.count=1且docs.deleted=149711的情况,核心原因是生成的document_id不唯一,导致每次数据导入时大量文档被重复写入同一个ID,旧版本文档被标记为删除,最终只剩最后一条写入的文档有效。此外,全量拉取数据配合update动作会加剧这个问题——即使ID唯一,全量重复同步也会生成大量墓碑文档。
分步解决方案
1. 验证并修复document_id的唯一性
首先确认拼接的复合键col1+col2是否真的唯一,且不存在空值:
- 临时修改Logstash配置,添加
filter和控制台输出,查看生成的document_id:
input { jdbc { jdbc_driver_library => "<path>/ojdbc10.jar" jdbc_driver_class => "Java::oracle.jdbc.driver.OracleDriver" jdbc_connection_string => "connection string" jdbc_user => <username> jdbc_password => <pwd> statement => "SELECT * FROM my_table_name" } } filter { mutate { add_field => { "generated_doc_id" => "%{col1}-%{col2}" } } } output { stdout { codec => rubydebug } elasticsearch { hosts => ["http://localhost:9200"] index => "test4" action => "update" doc_as_upsert => true document_id => "%{col1}-%{col2}" } }
启动Logstash后观察控制台输出的generated_doc_id:
- 若大量ID重复:检查Oracle表中
col1/col2是否存在空值,或复合键本身不满足唯一约束。 - 若存在空值:用
ruby脚本为字段设置默认值,避免ID拼接为空或重复:
document_id => "%{ruby: event.get('col1') || 'empty_col1'}-%{ruby: event.get('col2') || 'empty_col2'}"
2. 配置增量同步,避免全量重复写入
全量拉取配合update动作会不断生成旧文档的墓碑记录,必须启用JDBC输入的增量同步:
- 假设Oracle表有记录更新时间的字段(如
update_time,类型为timestamp),修改JDBC输入配置:
jdbc { jdbc_driver_library => "<path>/ojdbc10.jar" jdbc_driver_class => "Java::oracle.jdbc.driver.OracleDriver" jdbc_connection_string => "connection string" jdbc_user => <username> jdbc_password => <pwd> statement => "SELECT * FROM my_table_name WHERE update_time > :sql_last_value" use_column_value => true tracking_column => "update_time" tracking_column_type => "timestamp" last_run_metadata_path => "/var/log/logstash/oracle_last_run.txt" clean_run => false schedule => "*/5 * * * *" }
- 若表中无更新时间字段,可使用复合键配合序列追踪,或启用分页拉取,核心是只处理变化的数据。
3. 清理索引中的墓碑文档
对于已生成大量docs.deleted的test4索引,可在业务低峰期执行强制合并清理:
curl -XPOST "http://<ip>:9200/test4/_forcemerge?only_expunge_deletes=true&max_num_segments=1"
最终优化后的完整配置示例
input { jdbc { jdbc_driver_library => "<path>/ojdbc10.jar" jdbc_driver_class => "Java::oracle.jdbc.driver.OracleDriver" jdbc_connection_string => "connection string" jdbc_user => <username> jdbc_password => <pwd> statement => "SELECT * FROM my_table_name WHERE update_time > :sql_last_value" use_column_value => true tracking_column => "update_time" tracking_column_type => "timestamp" last_run_metadata_path => "/var/log/logstash/oracle_last_run.txt" clean_run => false schedule => "*/5 * * * *" } } output { elasticsearch { hosts => ["http://localhost:9200"] index => "test4" action => "update" doc_as_upsert => true document_id => "%{ruby: event.get('col1') || 'default_col1'}-%{ruby: event.get('col2') || 'default_col2'}" } }
内容的提问来源于stack exchange,提问作者Geo Thomas
相关产品推荐
相关产品推荐

