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

AWS Firehose向ElasticSearch发送数据存为字符串而非JSON对象问题排查

问题描述

我正在使用AWS Kinesis Firehose和ElasticSearch构建日志摄入服务,通过Node.js API将日志发送至Firehose。

发送的日志结构

{"level": "error","message": "Failed to connect to DB","resourceId": "server-1234","timestamp": "2023-09-15T08:00:00Z","traceId": "abc-xyz-123","spanId": "span-456","commit": "5e5342f","metadata": {"parentResourceId": "server-0987"}}

处理流程

由于Firehose要求接收字符串,我先通过JSON.stringify将日志对象转为字符串,再在Firehose传输流中使用Lambda函数将其转换回JSON对象,Lambda函数代码如下:

exports.handler = async (event) => {
const output = event.records.map((record) => {
    // Decoding the base64 data
    const payload = Buffer.from(record.data, 'base64').toString('utf8');
    let parsedData;

    try {
        parsedData = JSON.parse(payload);
    } catch (error) {
        console.error('Error parsing JSON:', error);
        return {
            recordId: record.recordId,
            result: 'ProcessingFailed',
            data: record.data,
        };
    }

    // Checking if 'message' exists and is a string that can be parsed as JSON
    if (typeof parsedData.message === 'string') {
        try {
            const messageData = JSON.parse(parsedData.message);
            parsedData = { ...parsedData, ...messageData };
            delete parsedData.message; 
        } catch (error) {
            console.error('Error parsing message JSON:', error);
            
        }
    }
    console.log(parsedData)
    // Re-encode the transformed data back to base64
    const outputData = Buffer.from(JSON.stringify(parsedData)).toString('base64');
    return {
        recordId: record.recordId,
        result: 'Ok',
        data: outputData,
    };
});

return { records: output };
};

但完成上述操作后,ElasticSearch索引中的数据仍以字符串形式存储。若直接将日志发送至ElasticSearch则会以正常JSON对象存储。请问为何会出现这种情况?是否无法通过Firehose发送JSON对象,或是我遗漏了某些配置?

索引映射

{  
  "mappings": {    
    "properties": {      
      "@timestamp": {        
        "type": "date"      
      },      
      "aws": {        
        "properties": {          
          "firehose": {            
            "properties": {              
              "arn": {                
                "type": "text",                
                "fields": {                  
                  "keyword": {                    
                    "type": "keyword",                    
                    "ignore_above": 256                  
                  }                
                }              
              },              
              "parameters": {                
                "properties": {                  
                  "es_datastream_name": {                    
                    "type": "text",                    
                    "fields": {                      
                      "keyword": {                        
                        "type": "keyword",                        
                        "ignore_above": 256                      
                      }                    
                    }                  
                  }                
                }              
              },              
              "request_id": {                
                "type": "text",                
                "fields": {                  
                  "keyword": {                    
                    "type": "keyword",                    
                    "ignore_above": 256                  
                  }                
                }              
              }            
            }          
          },          
          "kinesis": {            
            "properties": {              
              "name": {                
                "type": "text",                
                "fields": {                  
                  "keyword": {                    
                    "type": "keyword",                    
                    "ignore_above": 256                  
                  }                
                }              
              },              
              "type": {                
                "type": "text",                
                "fields": {                  
                  "keyword": {                    
                    "type": "keyword",                    
                    "ignore_above": 256                  
                  }                
                }              
              }            
            }          
          }        
        }      
      },      
      "cloud": {        
        "properties": {          
          "account": {            
            "properties": {              
              "id": {                
                "type": "text",                
                "fields": {                  
                  "keyword": {                    
                    "type": "keyword",                    
                    "ignore_above": 256                  
                  }                
                }              
              }            
            }          
          },          
          "provider": {            
            "type": "text",            
            "fields": {              
              "keyword": {                
                "type": "keyword",                
                "ignore_above": 256              
              }            
            }          
          },          
          "region": {            
            "type": "text",            
            "fields": {              
              "keyword": {                
                "type": "keyword",                
                "ignore_above": 256              
              }            
            }          
          }        
        }      
      },      
      "commit": {        
        "type": "keyword"      
      },      
      "level": {        
        "type": "keyword"      
      },      
      "message": {        
        "type": "text",        
        "fielddata": true      
      },      
      "metadata": {        
        "type": "nested",        
        "properties": {          
          "parentResourceId": {            
            "type": "keyword"          
          }        
        }      
      },      
      "resourceId": {        
        "type": "keyword"      
      },      
      "spanId": {        
        "type": "keyword"      
      },      
      "timestamp": {        
        "type": "date"      
      },      
      "traceId": {        
        "type": "keyword"      
      }    
  }  
}

问题原因及解决方法

核心原因

Firehose向Elasticsearch写入时,默认会将每条记录包裹在data字段中,导致Lambda处理后的JSON被当作字符串存放在该字段下,而非直接解析为索引的顶层字段。另外,若Firehose的Elasticsearch目标未正确配置数据格式,也会导致JSON无法被识别,最终以字符串形式存储。

解决步骤

1. 调整Firehose Elasticsearch目标配置

  • 打开Firehose传输流的Elasticsearch目标设置
  • 确认记录格式设置为JSON,开启JSON解析选项(若存在)
  • 检查是否启用了直接PUT模式,该模式会让Firehose直接将记录写入Elasticsearch,而非包裹在data字段中

2. 确认Lambda输出格式正确性

你的Lambda代码逻辑本身没问题,将处理后的JSON对象转为字符串再base64编码符合Firehose要求。但需确保:

  • 每条记录的data是单行JSON字符串的base64编码(避免换行,防止Firehose解析错误)
  • Lambda日志中parsedData是正确的JSON对象,无格式错误

3. 排查Elasticsearch文档结构

查询Elasticsearch中的文档,若发现所有日志字段都嵌套在data字段下,说明Firehose的默认包裹逻辑生效。此时必须调整Firehose配置,关闭自动包裹,或开启JSON直接解析。

4. 备选:临时适配索引映射(不推荐)

若无法修改Firehose配置,可临时修改索引映射,添加data字段为object类型,包含所有日志字段:

{
  "mappings": {
    "properties": {
      "data": {
        "properties": {
          "level": {"type": "keyword"},
          "message": {"type": "text"},
          "resourceId": {"type": "keyword"},
          "timestamp": {"type": "date"},
          "traceId": {"type": "keyword"},
          "spanId": {"type": "keyword"},
          "commit": {"type": "keyword"},
          "metadata": {
            "type": "nested",
            "properties": {
              "parentResourceId": {"type": "keyword"}
            }
          }
        }
      },
      // 保留原有aws、cloud等字段映射
      "aws": { /* ... */ },
      "cloud": { /* ... */ }
    }
  }
}

但此方案会导致字段嵌套,增加查询复杂度,仅作为临时过渡使用。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 21:30:06