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,提问作者מתן שוורץ

