通过Logstash将多份CSV文件导入带嵌套字段的单个索引
问题:Logstash多CSV关联合并为带嵌套结构的索引
我想用Logstash把Users、Scores、Messages三个CSV文件加载到同一个Elasticsearch索引里,所有文件都有userId字段用来关联数据。目标是生成users-index,其中用户CSV的内容作为普通字段,分数和消息数据做成嵌套字段(一个用户对应多条消息和分数)。我不确定自己的合并方式对不对,以下是我尝试的配置:
input { file { path => "C:/resources/files/users.csv" start_position => "beginning" sincedb_path => "NUL" } file { path => "C:/resources/files/scores.csv" start_position => "beginning" sincedb_path => "NUL" } file { path => "C:/resources/files/messages.csv" start_position => "beginning" sincedb_path => "NUL" } } filter { if [log][file][path] == "C:/resources/files/users.csv" { csv { separator => "," columns => ["userId", "username", "email"] } mutate { remove_field => ["[event][original]", "[log][file][path]", "[log][file][path][keyword]", "[message]", "[message][keyword]"] } } if [log][file][path] == "C:/resources/files/scores.csv" { csv{ separator => "," columns => ["userId", "field", "score"] } translate { destination => "[@metadata][scores]" dictionary_path => "C:/resources/files/scores.csv" field => "userId" } dissect { mapping => { "[@metadata][scores]" => "%{field};%{score}" } } } if [log][file][path] == "C:/resources/files/messages.csv" { csv { separator => "," columns => ["userId", "message", "tag"] } translate { destination => "[@metadata][messages]" dictionary_path => "C:/resources/files/messages.csv" field => "userId" } dissect { mapping => { "[@metadata][messages]" => "%{message};%{tag}" } } } } output { elasticsearch{ action => "create" hosts => "localhost:9200" index => "users-index" } }
请问有没有办法实现这个需求?
解决方案
你当前的配置没法实现需求,核心问题在于:
translate插件只能做单键值匹配,没法收集同一个userId对应的多条分数/消息记录- 现在的逻辑会把每个CSV的每行都单独输出成ES文档,没法把同一个用户的多条分数、消息合并到一个文档里的嵌套数组中
要实现这个需求,必须用Logstash的aggregate插件来做分组聚合,它可以根据userId把相关的用户、分数、消息记录聚合成一个完整的文档。下面是完整的可行配置:
input { file { path => "C:/resources/files/users.csv" start_position => "beginning" sincedb_path => "NUL" type => "users" # 给用户文件加类型标签,方便后续判断 } file { path => "C:/resources/files/scores.csv" start_position => "beginning" sincedb_path => "NUL" type => "scores" # 分数文件标签 } file { path => "C:/resources/files/messages.csv" start_position => "beginning" sincedb_path => "NUL" type => "messages" # 消息文件标签 } } filter { # 解析用户CSV if [type] == "users" { csv { separator => "," columns => ["userId", "username", "email"] } mutate { remove_field => ["event", "log", "message"] # 清理无关字段 add_field => { "[@metadata][is_user]" => "true" } } # 初始化聚合:以userId为分组ID,创建用户基础文档 aggregate { task_id => "%{userId}" code => " map['userId'] = event.get('userId') map['username'] = event.get('username') map['email'] = event.get('email') map['scores'] = [] unless map['scores'] map['messages'] = [] unless map['messages'] event.cancel() # 暂时取消输出,等聚合完成再输出 " } } # 解析分数CSV if [type] == "scores" { csv { separator => "," columns => ["userId", "field", "score"] } mutate { remove_field => ["event", "log", "message"] } # 将分数添加到对应userId的聚合数组中 aggregate { task_id => "%{userId}" code => " score = { 'field' => event.get('field'), 'score' => event.get('score') } map['scores'] << score event.cancel() " } } # 解析消息CSV if [type] == "messages" { csv { separator => "," columns => ["userId", "message_content", "tag"] # 避免和默认message字段冲突 } mutate { remove_field => ["event", "log", "message"] } # 将消息添加到对应userId的聚合数组中 aggregate { task_id => "%{userId}" code => " message = { 'content' => event.get('message_content'), 'tag' => event.get('tag') } map['messages'] << message event.cancel() " } } # 聚合完成后输出完整文档 aggregate { task_id => "%{userId}" timeout => 30 # 设置超时,确保所有关联记录都被聚合 code => " event.set('userId', map['userId']) event.set('username', map['username']) event.set('email', map['email']) event.set('scores', map['scores']) event.set('messages', map['messages']) " } } output { # 只输出聚合后的完整用户文档 if [userId] { elasticsearch { action => "index" hosts => "localhost:9200" index => "users-index" document_id => "%{userId}" # 用userId做文档ID,避免重复 } # 可选:输出到控制台调试 stdout { codec => rubydebug } } }
关键说明:
- 类型标签:给每个file输入加
type字段,比判断文件路径更简洁可靠 - aggregate插件:
- 以
userId作为task_id,实现按用户分组 - 处理用户记录时初始化聚合结构,包含基础字段和空的分数、消息数组
- 处理分数/消息时,将单条记录转为哈希,添加到对应用户的数组中
- 设置
timeout确保所有相关记录都被聚合后再输出完整文档
- 以
- 字段冲突处理:把消息CSV里的
message字段改成message_content,避免和Logstash默认的message字段冲突 - 文档ID:用
userId作为ES文档ID,确保同一个用户只会生成一条文档,不会重复
内容的提问来源于stack exchange,提问作者Ivan Andreev
相关产品推荐
相关产品推荐

