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

Elastic Stack-Logstash JSON解析错误排查求助

问题描述

基于Elastic Stack构建日志服务,通过TCP接收微服务日志,Logstash根据日志内容输出到HTTP或Elasticsearch(或两者)。处理大日志时触发JSON解析错误:

JSON parse error, original data now in message field {:message=>"Unexpected end-of-input: was expecting closing quote for a string value\n at [Source: (StringReader); line: 1, column: 9560693]", :exception=>LogStash::Json::ParserError}

当前logstash.conf配置如下:

input{
    tcp {
        port => 5000
        codec => json
    }
}
filter{
    if [system] == 'nostro' {
        if [event] == 'state_snapshot' {
            json {
                source => "state"
                target => "current_state"
            }
            mutate { remove_field => ["state"] }

            # Separate state fields
            ruby {
                code => '
                    current_state = event.get("current_state")
                    if current_state.is_a?(Hash)
                        current_state.each do |key, value|
                            if value.is_a?(Array)
                                parsed_arr = value.map { |item| JSON.parse(item) rescue item}
                                event.set("[current_state][#{key}]", parsed_arr)
                            elsif key.include?("orders") && value.is_a?(Hash)
                                parsed_orders = value.transform_values { |json_string| JSON.parse(json_string) rescue json_string }
                                event.set("[current_state][#{key}]", parsed_orders)
                            end
                        end
                    end
                '
            }

            mutate { rename => { "current_state" => "state" } }
        }
    }
}
output{
    # Condition the output index on the incoming logs 'system' field
    if [system] == 'nostro' {
        # Send to telegram if log is error
        if [log_level] in ['ERROR', 'CRITICAL'] {
            http {
                url => "http://host.docker.internal:8000/telegram-notification"
                http_method => "post"
                headers => {
                    "Content-Type" => "application/json"
                }
                message => ""
                automatic_retries => 1
                retry_failed => false
            }
        }
        if [event] == 'state_snapshot' {
            http {
                url => "http://host.docker.internal:8000/log-to-s3"
                http_method => "post"
                headers => {
                    "Content-Type" => "application/json"
                }
                message => ""
                automatic_retries => 2
                retry_failed => false
            }
        }
        elasticsearch {
            hosts => ["http://elasticsearch:9200"]
            index => "%{system}-%{service}-%{+YYYY.MM.dd}"
        }
    }
    else if [system] == 'riskmanagement' {
        elasticsearch {
            hosts => ["http://elasticsearch:9200"]
            index => "riskmanagement-%{+YYYY.MM.dd}"
        }
    }
    else {
        elasticsearch {
                hosts => ["http://elasticsearch:9200"]
                index => "other_index"
            }
    }
    stdout{codec => rubydebug}
}

已知测试情况:仅用HTTP输出该大日志时正常;将日志缩短为同格式(保留部分数组、哈希字段值)也可正常解析。

解决办法

1. 调整TCP输入的缓冲区与帧处理

TCP传输大日志时易出现数据截断,导致JSON不完整。修改TCP输入配置,增大缓冲区并启用按行解析:

input{
    tcp {
        port => 5000
        # 替换为json_lines codec,按行解析单条JSON日志
        codec => json_lines
        # 增大接收缓冲区至16MB,适配大日志(可根据实际大小调整)
        buffer_size => 16777216
    }
}

json_lines会以换行符为分隔符解析日志,避免大日志被截断导致的不完整问题。

2. 给JSON解析添加容错与调试机制

在filter阶段配置JSON解析的容错逻辑,同时保留原始数据便于排查问题:

filter{
    if [system] == 'nostro' {
        if [event] == 'state_snapshot' {
            json {
                source => "state"
                target => "current_state"
                # 解析失败时不中断流程,保留原始字段
                skip_on_invalid_json => true
                # 给解析失败的事件添加标签,方便后续处理
                add_tag => ["json_parse_failed"]
            }
            # 针对解析失败的事件,记录原始state内容到单独字段
            if "json_parse_failed" in [tags] {
                mutate { add_field => ["raw_state", "%{state}"] }
            }
            mutate { remove_field => ["state"] }

            # 后续ruby逻辑保持不变...
        }
    }
}

这样即使解析失败,也能留存原始数据定位问题点。

3. 优化Ruby代码的JSON解析容错

当前Ruby代码的rescue逻辑过于笼统,无法定位具体哪个元素解析失败。修改为带错误记录的解析逻辑:

ruby {
    code => '
        current_state = event.get("current_state")
        if current_state.is_a?(Hash)
            current_state.each do |key, value|
                if value.is_a?(Array)
                    parsed_arr = []
                    value.each_with_index do |item, idx|
                        begin
                            parsed_item = JSON.parse(item)
                            parsed_arr << parsed_item
                        rescue JSON::ParserError => e
                            # 记录解析失败的数组元素索引和内容
                            event.set("[parse_errors][#{key}_array_idx_#{idx}]", item)
                            parsed_arr << item
                        end
                    end
                    event.set("[current_state][#{key}]", parsed_arr)
                elsif key.include?("orders") && value.is_a?(Hash)
                    parsed_orders = {}
                    value.each do |order_key, json_string|
                        begin
                            parsed_order = JSON.parse(json_string)
                            parsed_orders[order_key] = parsed_order
                        rescue JSON::ParserError => e
                            # 记录解析失败的订单字段内容
                            event.set("[parse_errors][#{key}_field_#{order_key}]", json_string)
                            parsed_orders[order_key] = json_string
                        end
                    end
                    event.set("[current_state][#{key}]", parsed_orders)
                end
            end
        end
    '
}

该逻辑会精准记录哪个元素解析失败,快速定位大日志中的问题片段。

4. 调整Elasticsearch的字段大小限制

Elasticsearch默认对字符串字段有ignore_above限制(超过1024KB会被截断),导致后续解析异常。可通过索引模板调整该参数:

{
  "index_patterns": ["nostro-*"],
  "mappings": {
    "properties": {
      "state": {
        "type": "object",
        "dynamic": true,
        "properties": {
          "large_content_field": {
            "type": "text",
            "ignore_above": 2097152 // 调整为2MB,根据实际需求设置
          }
        }
      }
    }
  }
}

也可在Logstash的Elasticsearch输出中开启模板管理,自定义字段限制规则。

5. 拆分大日志为多条事件

若日志体积过大,可在filter阶段拆分大数组/哈希为多条小事件,降低单条事件的处理压力:

ruby {
    code => '
        current_state = event.get("current_state")
        if current_state.is_a?(Hash) && current_state.key?("orders")
            orders_list = current_state.delete("orders")
            orders_list.each do |order_item|
                clone_event = event.clone
                clone_event.set("single_order", order_item)
                clone_event.tag("split_order")
                yield clone_event
            end
        end
    '
}

拆分后单条事件体积减小,可避免解析和存储时的大小限制问题。


内容的提问来源于stack exchange,提问作者מתן שוורץ

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 11:18:10