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

修复Kinesis Firehose转换Lambda的CloudWatch日志JSON格式错误

问题:CloudWatch日志转存S3后JSON格式异常及Lambda修改报错

我已搭建通过Kinesis Firehose将CloudWatch日志转存至S3的架构,但S3中生成的日志文件并非标准JSON数组格式——标准格式应为[{},{},{}],而Firehose生成的是{}{}{}。

我尝试修改「Process CloudWatch logs sent to Kinesis Firehose」Lambda蓝图来修正格式:

  1. 在transformLogEvent方法中修改返回内容为log_event['message'] + ',\n',给每个JSON对象添加逗号换行
  2. 在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数组,避免后续拼接出错,同时自带格式校验:

  1. 修改transformLogEvent方法:
def transformLogEvent(log_event):
    return log_event['message']
  1. 修改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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 15:18:09