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

如何在Logstash中使用map实现MySQL到Elasticsearch的数据合并

问题描述

目标是合并MySQL分片数据,数据表结构如下:

shardcommunityIdpostIdjson
BIGINTVAR64VAR32{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
  }
}
问题排查与修复方案

核心问题点

  1. 字段名称不匹配:数据表中存储JSON的字段名为json,但配置中错误使用event.get('data'),导致无法读取目标内容
  2. 聚合逻辑错误:当前仅将JSON对象存入数组,未实现合并;且push_previous_map_as_event设为false,聚合完成后不会生成新事件
  3. 缺少事件过滤逻辑:未处理未完成聚合的事件,也没有对超时后的聚合结果做输出控制

修复后的完整配置

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 01:30:56