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

添加复合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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 08:50:22