I am trying to stream my cloudwatch logs to AWS Opensearch(ElasticSearch) by creating a kinesis firehose subscription filter. I want to understand does Kinesis data firehose supports bulk insert to the Elastic Search via bulk API or it inserts single records ?
So basically, I have created a lambda transformer which aggregates the Cloudwatch Log messages and transforms it to the format required for the bulk insert API of elastic search, I have tested my data format by directly calling the bulk API of elasticsearch and it works fine, however when I try to do the same via firehose I get following error :
message : One or more records are malformed. Please ensure that each record is single valid JSON object and that it does not contain newlines
error code : OS.MalformedData
The transformed data :
{"index":{"_index":"test","_type":"/aws/lambda/HelloWorldLogProducer","_id":"36495659535505340260631192776509739326813505134549729280"}}
{"timeMillis":1636521972493,"thread":"main","level":"INFO","loggerName":"helloworld.App","message":"Inflating logs to increase size","endOfBatch":false,"loggerFqcn":"org.apache.logging.slf4j.Log4jLogger","contextMap":{"AWSRequestId":"332de390-42b3-44dd-937a-cc4089ff9510"},"threadId":1,"threadPriority":5,"jobId":"${ctx:request_id}","clientUUID":"${ctx:client_uuid}","clientUUIDHeader":"${ctx:client_uuid_header}"}
{"index":{"_index":"test","_type":"/aws/lambda/HelloWorldLogProducer","_id":"36495659535884452929006213369915846537448527280151396358"}}
{"timeMillis":1636521973189,"thread":"main","level":"INFO","loggerName":"helloworld.App","message":"{ \"message\": \"hello world\", \"location\": \"3.137.161.28\" }","endOfBatch":false,"loggerFqcn":"org.apache.logging.slf4j.Log4jLogger","contextMap":{"AWSRequestId":"332de390-42b3-44dd-937a-cc4089ff9510"},"threadId":1,"threadPriority":5,"jobId":"${ctx:request_id}","clientUUID":"${ctx:client_uuid}","clientUUIDHeader":"${ctx:client_uuid_header}"}
It seems that firehose does not supports bulk insert, like the cloud watch to opensearch subscription filter does, however I cannot define the index name in opensearch subscription filter.
Does anyone faced a similar problem before ?