求助:修正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
相关产品推荐
相关产品推荐

