如何在AWS Lambda中将转换后的Parquet文件写入S3?
问题描述
现有一个包含A、B两列的Parquet文件,A列数据类型为string,B列数据类型为float64。需要将B列数据类型修改为int64,由于无法直接修改Parquet文件,已在本地成功创建符合目标数据类型(A为string、B为int64)的新Parquet文件。现在需要在AWS Lambda中复现该逻辑,但在最后一步将转换后的Parquet文件写入S3时遇到问题。
现有Lambda代码如下:
import pyarrow as pa import pyarrow.parquet as pq import boto3 from io import BytesIO def lambda_handler(event, context): s3_bucket = 'a_bucket' s3_client = boto3.client('s3') key_name = 'testfolder/testfile.parquet' objSh1 = s3_client.get_object(Bucket=s3_bucket, Key=key_name) pq_raw = pq.read_table(source=BytesIO(objSh1['Body'].read())) df_raw = pq_raw.to_pandas() sc = pq_raw.schema for column, datatype in zip(sc.names, sc.types): print(f"{column}, -> {datatype}") print(df_raw.dtypes) test_pq_raw = pa.Table.from_pandas( df=df_raw, preserve_index=False ) sc_test = test_pq_raw.schema for column, datatype in zip(sc_test.names, sc_test.types): print(f"{column}, -> {datatype}") schema = { 'A':pa.string(), 'B':pa.int64() } fields = [pa.field(x, y) for x, y in schema.items()] new_schema = pa.schema(fields) table = pa.Table.from_pandas( df_raw, schema=new_schema, preserve_index=False )
本地环境中用于写入Parquet文件的代码:
pq.write_table(table, 'C:\\Users\\username1\\Desktop\\testfolder\\testoutput.parquet')
需要实现Lambda中向S3写入转换后Parquet文件的操作。
解决方案
Lambda中无法直接写入本地文件,需借助内存缓冲区(BytesIO)存储转换后的Parquet数据,再通过boto3上传到S3。在现有代码基础上添加以下步骤:
- 创建
BytesIO对象作为内存缓冲区 - 将转换后的表写入缓冲区
- 把缓冲区内容上传到指定S3路径
修改后的完整Lambda代码如下:
import pyarrow as pa import pyarrow.parquet as pq import boto3 from io import BytesIO def lambda_handler(event, context): s3_bucket = 'a_bucket' s3_client = boto3.client('s3') input_key = 'testfolder/testfile.parquet' output_key = 'testfolder/testoutput.parquet' # 定义输出文件的S3路径 # 读取S3中的原始Parquet文件 objSh1 = s3_client.get_object(Bucket=s3_bucket, Key=input_key) pq_raw = pq.read_table(source=BytesIO(objSh1['Body'].read())) df_raw = pq_raw.to_pandas() # 可选:打印原表结构用于调试 sc = pq_raw.schema for column, datatype in zip(sc.names, sc.types): print(f"{column}, -> {datatype}") print(df_raw.dtypes) # 定义目标Schema schema = { 'A': pa.string(), 'B': pa.int64() } fields = [pa.field(x, y) for x, y in schema.items()] new_schema = pa.schema(fields) # 转换数据类型 table = pa.Table.from_pandas( df_raw, schema=new_schema, preserve_index=False ) # 将转换后的表写入内存缓冲区 buffer = BytesIO() pq.write_table(table, buffer) buffer.seek(0) # 将文件指针移到缓冲区开头 # 上传缓冲区内容到S3 s3_client.put_object( Bucket=s3_bucket, Key=output_key, Body=buffer.getvalue() ) return { 'statusCode': 200, 'body': f"转换后的Parquet文件已成功写入S3路径:s3://{s3_bucket}/{output_key}" }
注意事项
- 确保Lambda执行角色拥有目标S3存储桶的
putObject权限 - 若B列float64数据包含小数,转换为int64时会自动截断小数(如
1.9转为1),如需四舍五入,可在转换前处理DataFrame:df_raw['B'] = df_raw['B'].round().astype('int64') - 内存缓冲区使用后必须调用
seek(0)重置指针,否则上传内容为空
内容的提问来源于stack exchange,提问作者Ludwig
相关产品推荐
相关产品推荐

