如何通过Kinesis Data Firehose转换CloudWatch日志写入ElasticSearch并解决批量插入报错
问题核心原因
首先明确:Kinesis Data Firehose对接AWS OpenSearch/Elasticsearch时,底层默认就是通过bulk API做批量写入的,不需要你自行拼接bulk API要求的双行元数据+文档格式,你当前的报错就是因为自行拼接了bulk格式导致的。
Firehose的处理逻辑是:
- 它要求Lambda转换器返回的每条记录都是无换行的单条合法JSON文档对象,不需要带
{"index": {...}}这类元数据行 - Firehose会自动为每批记录拼接符合bulk API要求的元数据头,再统一调用OpenSearch的bulk接口写入,你自行添加的元数据行反而会导致单条记录格式不符合要求,触发
OS.MalformedData报错。
自定义索引名称的实现方式
你需要的自定义索引能力不需要通过自行拼接bulk格式实现,直接通过Firehose的OpenSearch目标配置即可完成:
- 如果需要按记录内容动态指定索引:在Lambda转换生成的单条文档JSON中加入一个自定义字段(比如
index_name),然后在Firehose的OpenSearch目标配置的「索引名称」项填写${index_name},Firehose会自动读取每条记录的对应字段值作为写入的索引名 - 如果是固定前缀+时间后缀的索引规则:直接在索引名称配置里填写模板即可,比如
test-%{+YYYY.MM.dd},Firehose会自动按写入时间替换为对应日期后缀 - 如果你还需要自定义文档ID、类型等属性:可以分别在Firehose配置的「文档ID字段」「类型」项指定对应值即可,OpenSearch 7.x及以上版本建议将_type设置为
_doc即可。
调整步骤
- 修改Lambda转换逻辑,去掉你拼接的
{"index": {...}}元数据行,仅返回纯文档内容的单条无换行JSON - 按上述说明调整Firehose的OpenSearch目标配置,指定你需要的索引规则
- 重新推送数据测试即可正常写入。
内容的提问来源于stack exchange,提问作者Paras
相关产品推荐
相关产品推荐

