如何在Logstash中使用map实现MySQL到Elasticsearch的数据合并
问题描述
目标是合并MySQL分片数据,数据表结构如下:
| shard | communityId | postId | json |
|---|---|---|---|
| BIGINT | VAR64 | VAR32 | {valid json with nested encoded json} |
同一communityId和postId的分片存储部分JSON内容,例如:
1,cid,pid,{title:"123"}
和
1,cid,pid,{desc:"desc here"}
希望在Logstash中通过简单SELECT查询(不使用GROUP_CONCAT)完成合并,当前配置未成功,配置代码如下:
input { jdbc { jdbc_driver_library => "/usr/mysql-connector-j-utf8.jar" #valid jdbc_driver_class => "com.mysql.cj.jdbc.Driver" jdbc_connection_string => "jdbc:mysql://*/table" jdbc_user => "*" jdbc_password => "*" schedule => "* * * * *" # Schedule for querying data (adjust as needed) statement => "SELECT * FROM sharded" } } filter { # Ensure events are sorted by 'communityId' and 'postId' ruby { code => " event.set('[@metadata][sort_key]', [event.get('communityId'), event.get('postId')].join('-')) " } # Aggregate events based on 'communityId' and 'postId' aggregate { task_id => "%{communityId}-%{postId}" code => " map['communityId'] = event.get('communityId') map['postId'] = event.get('postId') map['combined_data'] ||= [] map['combined_data'].push(JSON.parse(event.get('data'))) event.cancel()" push_previous_map_as_event => false timeout => 3 # Adjust the timeout as needed (in seconds) timeout_tags => ['aggregated'] } # Use a conditional to filter out events without 'combined_data' } output { elasticsearch { hosts => ["elasticsearch-container:9200"] index => "collections" document_id => "%{[@metadata][sort_key]}" # Use lowercase 'postId' action => "update" doc_as_upsert => true } stdout { codec => rubydebug } }
问题排查与修复方案
核心问题点
- 字段名称不匹配:数据表中存储JSON的字段名为
json,但配置中错误使用event.get('data'),导致无法读取目标内容 - 聚合逻辑错误:当前仅将JSON对象存入数组,未实现合并;且
push_previous_map_as_event设为false,聚合完成后不会生成新事件 - 缺少事件过滤逻辑:未处理未完成聚合的事件,也没有对超时后的聚合结果做输出控制
修复后的完整配置
input { jdbc { jdbc_driver_library => "/usr/mysql-connector-j-utf8.jar" jdbc_driver_class => "com.mysql.cj.jdbc.Driver" jdbc_connection_string => "jdbc:mysql://*/table" jdbc_user => "*" jdbc_password => "*" schedule => "* * * * *" # 查询时按communityId和postId排序,确保同组事件连续处理 statement => "SELECT * FROM sharded ORDER BY communityId, postId" } } filter { ruby { code => " event.set('[@metadata][sort_key]', [event.get('communityId'), event.get('postId')].join('-')) " } aggregate { task_id => "%{communityId}-%{postId}" code => " # 初始化合并后的JSON对象 map['combined_json'] ||= {} # 读取当前事件的JSON内容并合并到map中 current_json = JSON.parse(event.get('json')) map['combined_json'].merge!(current_json) # 保留基础标识字段 map['communityId'] = event.get('communityId') map['postId'] = event.get('postId') # 取消原事件,避免单独输出分片数据 event.cancel() " # 超时后将聚合结果作为新事件输出 push_previous_map_as_event => true # 根据数据量调整超时时间,确保同组事件都被处理 timeout => 10 timeout_tags => ['aggregated'] } # 只保留聚合完成的事件,丢弃未完成的无效事件 if 'aggregated' not in [tags] { drop {} } # 可选:将合并后的JSON展开到事件根节点(如需直接写入ES根字段可启用) # ruby { # code => " # combined = event.get('combined_json') # combined.each { |k, v| event.set(k, v) } # event.remove('combined_json') # " # } } output { elasticsearch { hosts => ["elasticsearch-container:9200"] index => "collections" document_id => "%{[@metadata][sort_key]}" action => "update" doc_as_upsert => true } stdout { codec => rubydebug } }
关键优化说明
- 字段修正:将
event.get('data')改为event.get('json'),匹配数据表字段名 - 聚合逻辑优化:用
merge!方法将多个分片JSON合并为单个完整对象,替代数组存储 - 事件输出配置:将
push_previous_map_as_event设为true,确保超时后输出合并后的完整事件 - 数据排序:SQL查询加入
ORDER BY communityId, postId,让同组事件连续进入聚合插件,减少超时等待时间 - 事件过滤:添加条件判断丢弃未完成聚合的事件,避免无效分片数据流入输出端
内容的提问来源于stack exchange,提问作者user655355
相关产品推荐
相关产品推荐

