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

