修复Kinesis Firehose转换Lambda的CloudWatch日志JSON格式错误
问题:CloudWatch日志转存S3后JSON格式异常及Lambda修改报错
我已搭建通过Kinesis Firehose将CloudWatch日志转存至S3的架构,但S3中生成的日志文件并非标准JSON数组格式——标准格式应为[{},{},{}],而Firehose生成的是{}{}{}。
我尝试修改「Process CloudWatch logs sent to Kinesis Firehose」Lambda蓝图来修正格式:
- 在
transformLogEvent方法中修改返回内容为log_event['message'] + ',\n',给每个JSON对象添加逗号换行 - 在
lambda_handler末尾添加代码,给首尾记录加上方括号:
last = len(event['records'])-1 start = '['+base64.b64decode(records[0]['data']).decode('utf-8') end = base64.b64decode(records[last]['data']).decode('utf-8')+']' records[0]['data'] = base64.b64encode(start.encode('utf-8')) records[last]['data'] = base64.b64encode(end.encode('utf-8'))
但测试Lambda时触发KeyError: 'data'错误,错误栈如下:
{ "errorMessage": "'data'", "errorType": "KeyError", "requestId": "04ac151c-c429-484d-813f-ffd5d65286e2", "stackTrace": [ " File \"/var/task/lambda_function.py\", line 270, in lambda_handler\n start = '['+base64.b64decode(records[0]['data']).decode('utf-8')\n" ] }
相关Lambda蓝图代码如下:
import base64 import json import gzip import boto3 def transformLogEvent(log_event): """Transform each log event. The default implementation below just extracts the message and appends a newline to it. Args: log_event (dict): The original log event. Structure is {"id": str, "timestamp": long, "message": str} Returns: str: The transformed log event. """ return log_event['message'] + ',\n' def processRecords(records): for r in records: data = loadJsonGzipBase64(r['data']) recId = r['recordId'] # CONTROL_MESSAGE are sent by CWL to check if the subscription is reachable. # They do not contain actual data. if data['messageType'] == 'CONTROL_MESSAGE': yield { 'result': 'Dropped', 'recordId': recId } elif data['messageType'] == 'DATA_MESSAGE': joinedData = ''.join([transformLogEvent(e) for e in data['logEvents']]) dataBytes = joinedData.encode("utf-8") encodedData = base64.b64encode(dataBytes).decode('utf-8') yield { 'data': encodedData, 'result': 'Ok', 'recordId': recId } else: yield { 'result': 'ProcessingFailed', 'recordId': recId } def loadJsonGzipBase64(base64Data): return json.loads(gzip.decompress(base64.b64decode(base64Data))) def lambda_handler(event, context): isSas = 'sourceKinesisStreamArn' in event streamARN = event['sourceKinesisStreamArn'] if isSas else event['deliveryStreamArn'] region = streamARN.split(':')[3] streamName = streamARN.split('/')[1] records = list(processRecords(event['records'])) projectedSize = 0 recordListsToReingest = [] for idx, rec in enumerate(records): originalRecord = event['records'][idx] if rec['result'] != 'Ok': continue # If a single record is too large after processing, split the original CWL data into two, each containing half # the log events, and re-ingest both of them (note that it is the original data that is re-ingested, not the # processed data). If it's not possible to split because there is only one log event, then mark the record as # ProcessingFailed, which sends it to error output. if len(rec['data']) > 6000000: cwlRecord = loadJsonGzipBase64(originalRecord['data']) if len(cwlRecord['logEvents']) > 1: rec['result'] = 'Dropped' recordListsToReingest.append( [createReingestionRecord(isSas, originalRecord, data) for data in splitCWLRecord(cwlRecord)]) else: rec['result'] = 'ProcessingFailed' print(('Record %s contains only one log event but is still too large after processing (%d bytes), ' + 'marking it as %s') % (rec['recordId'], len(rec['data']), rec['result'])) del rec['data'] else: projectedSize += len(rec['data']) + len(rec['recordId']) # 6000000 instead of 6291456 to leave ample headroom for the stuff we didn't account for if projectedSize > 6000000: recordListsToReingest.append([createReingestionRecord(isSas, originalRecord)]) del rec['data'] rec['result'] = 'Dropped' # call putRecordBatch/putRecords for each group of up to 500 records to be re-ingested if recordListsToReingest: recordsReingestedSoFar = 0 client = boto3.client('kinesis' if isSas else 'firehose', region_name=region) maxBatchSize = 500 flattenedList = [r for sublist in recordListsToReingest for r in sublist] for i in range(0, len(flattenedList), maxBatchSize): recordBatch = flattenedList[i:i + maxBatchSize] # last argument is maxAttempts args = [streamName, recordBatch, client, 0, 20] if isSas: putRecordsToKinesisStream(*args) else: putRecordsToFirehoseStream(*args) recordsReingestedSoFar += len(recordBatch) print('Reingested %d/%d' % (recordsReingestedSoFar, len(flattenedList))) print('%d input records, %d returned as Ok or ProcessingFailed, %d split and re-ingested, %d re-ingested as-is' % ( len(event['records']), len([r for r in records if r['result'] != 'Dropped']), len([l for l in recordListsToReingest if len(l) > 1]), len([l for l in recordListsToReingest if len(l) == 1]))) # encapsulate in square brackets for proper JSON formatting last = len(event['records'])-1 start = '['+base64.b64decode(records[0]['data']).decode('utf-8') end = base64.b64decode(records[last]['data']).decode('utf-8')+']' records[0]['data'] = base64.b64encode(start.encode('utf-8')) records[last]['data'] = base64.b64encode(end.encode('utf-8')) return {'records': records}
错误原因分析
KeyError: 'data'的核心原因是并非所有records元素都包含data字段:
- 当记录是
CONTROL_MESSAGE时,processRecords会返回仅含result: Dropped和recordId的对象,无data字段 - 当处理后的记录过大需要重 ingest 时,代码会执行
del rec['data']并将result设为Dropped,此时该记录也无data字段 - 你的代码直接访问
records[0]['data'],若第一条记录是上述两种情况之一,就会触发KeyError。
修复方案
方案1:安全处理首尾记录
先筛选出所有result为Ok的有效记录,再对这些记录进行格式修正,同时处理最后一条记录末尾多余的逗号:
# 替换原lambda_handler末尾的格式修正代码 valid_records = [rec for rec in records if rec.get('result') == 'Ok'] if valid_records: # 处理第一条记录,添加开头的[并移除末尾多余逗号 first_data = base64.b64decode(valid_records[0]['data']).decode('utf-8') if first_data.endswith(',\n'): first_data = first_data[:-2] valid_records[0]['data'] = base64.b64encode(f'[{first_data}'.encode('utf-8')).decode('utf-8') # 处理最后一条记录,添加结尾的]并移除末尾多余逗号 last_data = base64.b64decode(valid_records[-1]['data']).decode('utf-8') if last_data.endswith(',\n'): last_data = last_data[:-2] valid_records[-1]['data'] = base64.b64encode(f'{last_data}]'.encode('utf-8')).decode('utf-8')
方案2:更优的JSON格式化方案(适配gzip压缩日志)
直接在处理DATA_MESSAGE时将日志事件转为标准JSON数组,避免后续拼接出错,同时自带格式校验:
- 修改
transformLogEvent方法:
def transformLogEvent(log_event): return log_event['message']
- 修改
processRecords中的DATA_MESSAGE处理逻辑:
elif data['messageType'] == 'DATA_MESSAGE': try: # 将每个日志消息解析为JSON对象,收集到列表中 log_entries = [] for e in data['logEvents']: try: log_entry = json.loads(e['message']) log_entries.append(log_entry) except json.JSONDecodeError: print(f"Invalid JSON in log event {e['id']}: {e['message']}") continue # 转为标准JSON数组字符串 json_output = json.dumps(log_entries) + '\n' dataBytes = json_output.encode("utf-8") encodedData = base64.b64encode(dataBytes).decode('utf-8') yield { 'data': encodedData, 'result': 'Ok', 'recordId': recId } except Exception as e: print(f"Error processing record {recId}: {str(e)}") yield { 'result': 'ProcessingFailed', 'recordId': recId }
该方案优势:
- 直接生成标准JSON数组,无需后续拼接首尾和处理逗号
- 自带JSON格式校验,能处理非法格式的日志条目
- 完美适配gzip压缩的CloudWatch日志(
loadJsonGzipBase64已处理解压和解析)
内容的提问来源于stack exchange,提问作者Brian G
相关产品推荐
相关产品推荐

