使用Logstash同步MSSQL到Elasticsearch时嵌套文档聚合异常求助
解决Logstash聚合MSSQL数据到Elasticsearch嵌套文档的异常问题
我帮你排查了下Logstash聚合配置里的问题,主要有几个关键调整点,修改后就能生成你想要的嵌套文档结构了:
问题分析与修正方案
1. 聚合推送逻辑错误
你当前设置了push_previous_map_as_event => true,这会导致每处理一条新事件就推送一次之前的聚合结果,产生大量重复且不完整的文档。应该把这个参数改为false,只在聚合任务完成(超时或所有事件处理完毕)时推送最终的聚合文档。
2. 字段大小写不匹配
你的Elasticsearch映射里嵌套字段是empid(小写d),但聚合代码里写的是'empId'(大写I),这会导致ES无法识别该字段,最终嵌套结构不符合预期。需要统一字段名的大小写。
3. 丢弃原始单条事件
原始的JDBC输入事件(每个部门+单个员工的组合)会被直接传递到输出,导致ES里既有聚合后的嵌套文档,又有原始的单条数据。需要添加drop过滤器,只保留聚合完成的事件。
4. 移除废弃的document_type参数
Elasticsearch 7.x及以后版本已经废弃了文档类型(document_type),保留这个参数会引发错误,需要删除该行。
修正后的Logstash配置
input { jdbc { jdbc_driver_library => "" jdbc_driver_class => "com.microsoft.sqlserver.jdbc.SQLServerDriver" jdbc_connection_string => "jdbc:sqlserver://EC2AMAZ-J90JR4A\SQLEXPRESS:1433;databaseName=xxxx;" jdbc_user => "xxxx" jdbc_password => "xxxx" statement => "SELECT De.id AS id,De.name AS deptname,Emp.id AS empid,Emp.emp_name AS empname FROM department De LEFT JOIN employee Emp ON De.id = Emp.dept_Id ORDER BY De.id" } } filter{ aggregate { task_id => "%{id}" code => " map['id'] = event['id'] map['deptname'] = event['deptname'] map['employee_details'] ||= [] # 修正字段名大小写,和ES映射保持一致 map['employee_details'] << {'empid' => event['empid'], 'empname' => event['empname'] } " # 关闭提前推送,只在任务完成时推送最终结果 push_previous_map_as_event => false # 确保所有任务都能被处理完成 flush_timeout => 10 timeout => 5 timeout_tags => ['aggregated'] } # 只保留聚合完成的事件,丢弃原始的单条数据 if 'aggregated' not in [tags] { drop {} } } output{ stdout{ codec => rubydebug } elasticsearch{ hosts => "https://d9bc7cbca5ec49ea96a6ea683f70caca.eastus2.azure.elastic-cloud.com:4567" user => "elastic" password => "****" index => "departments" action => "index" # 删除废弃的document_type参数 document_id => "%{id}" } }
额外说明
- 确保你的JDBC驱动包(
sqljdbc42.jar或对应版本)已经放在Logstash的jdbc_driver_library指定路径,或者Logstash的vendor/jars目录下,否则会出现驱动加载失败的错误。 - 运行Logstash后,你会在stdout里看到每个部门对应一条包含所有员工的嵌套文档,这样导入到Elasticsearch后就完全符合你期望的结构了。
内容的提问来源于stack exchange,提问作者Mohan vel
相关产品推荐
相关产品推荐

