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

通过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 }
    }
}

关键说明:

  1. 类型标签:给每个file输入加type字段,比判断文件路径更简洁可靠
  2. aggregate插件:
    • 以userId作为task_id,实现按用户分组
    • 处理用户记录时初始化聚合结构,包含基础字段和空的分数、消息数组
    • 处理分数/消息时,将单条记录转为哈希,添加到对应用户的数组中
    • 设置timeout确保所有相关记录都被聚合后再输出完整文档
  3. 字段冲突处理:把消息CSV里的message字段改成message_content,避免和Logstash默认的message字段冲突
  4. 文档ID:用userId作为ES文档ID,确保同一个用户只会生成一条文档,不会重复

内容的提问来源于stack exchange,提问作者Ivan Andreev

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 08:15:52