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

求助:修正AWS Lambda代码实现JSON转Parquet并直接返回

JSON转Parquet的AWS Lambda函数修正

问题分析

原代码存在以下关键问题:

  • 未定义processed_messages变量,遍历logEvents时未收集结构化数据
  • 未解析message字段中的键值对格式内容,无法生成正确的DataFrame
  • 未正确提取Parquet的字节数据,错误返回空字符串的base64编码
  • 包含未使用的boto3依赖

修正后的代码

import base64
import gzip
import json
import pandas as pd
import pyarrow as pa
import pyarrow.parquet as pq


def parse_log_message(message):
    # 提取message中键值对部分(跳过日志前缀)
    kv_part = message.split(' - ')[-1]
    # 分割键值对并转成字典
    kv_pairs = kv_part.split(',')
    result = {}
    for pair in kv_pairs:
        key, value = pair.split('=', 1)
        # 尝试转换数值类型
        try:
            if '.' in value:
                result[key] = float(value)
            else:
                result[key] = int(value)
        except ValueError:
            result[key] = value.strip()
    return result


def lambda_handler(event, context):
    output = []
    for record in event['records']:
        # 解码并解压输入数据
        compressed_data = base64.b64decode(record['data'])
        data = json.loads(gzip.decompress(compressed_data))
        
        processed_messages = []
        # 处理每条日志事件
        for log_event in data['logEvents']:
            # 解析message字段的结构化数据
            parsed_data = parse_log_message(log_event['message'])
            # 保留原日志的id和timestamp
            parsed_data['log_id'] = log_event['id']
            parsed_data['log_timestamp'] = log_event['timestamp']
            processed_messages.append(parsed_data)
        
        # 生成DataFrame
        df = pd.DataFrame(processed_messages)
        # 转换为Arrow Table
        table = pa.Table.from_pandas(df)
        
        # 内存中写入Parquet
        parquet_stream = pa.BufferOutputStream()
        pq.write_table(table, parquet_stream)
        parquet_bytes = parquet_stream.getvalue().to_pybytes()
        
        # 生成输出记录:base64编码Parquet字节
        output_record = {
            'recordId': record['recordId'],
            'result': 'Ok',
            'data': base64.b64encode(parquet_bytes).decode('utf-8')
        }
        output.append(output_record)
    
    return {'records': output}

关键修正说明

  • 新增parse_log_message函数:解析message字段中的键值对格式,将非结构化日志内容转为结构化字典,并自动转换数值类型(int/float)
  • 正确收集processed_messages:遍历logEvents时,将解析后的结构化数据存入列表,用于生成DataFrame
  • 提取Parquet字节数据:使用parquet_stream.getvalue().to_pybytes()获取内存中的Parquet字节内容
  • 正确生成输出数据:将Parquet字节进行base64编码后返回,符合Lambda(尤其是Firehose转换Lambda)的输出格式
  • 移除未使用的boto3依赖,减少包体积

内容的提问来源于stack exchange,提问作者Divya

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 12:42:41