在AWS Lambda中无需Spark,用Python/TypeScript写入S3 Parquet格式
在AWS Lambda中无Spark写入Parquet到S3(支持追加)
下面分别提供Python和TypeScript的实现方案,同时针对Parquet追加的特性给出合理的实现方式(注:Parquet单个文件不适合频繁追加,通常推荐向同一S3前缀下写入新的Parquet文件来实现逻辑上的追加,若需修改单个文件则需额外处理)。
Python 实现
依赖准备
需要安装 pyarrow 和 s3fs 库,Lambda中可以通过层或者打包成zip包部署(注意打包时要适配Lambda的运行环境,比如用Amazon Linux的镜像打包)。
代码示例
import pyarrow as pa import pyarrow.dataset as ds import s3fs import pandas as pd def lambda_handler(event, context): # 1. 准备待写入的数据(示例数据) data = { "id": [1, 2, 3], "name": ["Alice", "Bob", "Charlie"], "age": [25, 30, 35] } df = pd.DataFrame(data) table = pa.Table.from_pandas(df) # 2. 配置S3路径和文件系统 s3_path = "s3://your-bucket-name/parquet-data/" fs = s3fs.S3FileSystem() # 3. 写入/追加Parquet数据 # mode='append'会自动在指定前缀下添加新的Parquet文件,保持schema一致 ds.write_to_dataset( table, base_dir=s3_path, format="parquet", mode="append", filesystem=fs ) return { "statusCode": 200, "body": "数据成功写入/追加至S3 Parquet" }
关键说明
- 此方案通过
pyarrow.dataset实现逻辑追加:每次运行会在指定S3前缀下生成新的Parquet文件,查询时可通过读取该前缀下所有文件合并数据,这是Parquet生态中推荐的追加方式。 - Lambda执行角色需要配置S3的
PutObject和ListBucket权限。 - 确保每次追加的schema与已有数据一致,否则会写入失败。
TypeScript 实现
依赖准备
需要安装 @aws-sdk/client-s3 和 parquetjs 库,打包时需包含这些依赖(注意Lambda的Node.js版本兼容性)。
代码示例(逻辑追加:写入新文件到前缀)
import { S3Client, PutObjectCommand } from "@aws-sdk/client-s3"; import { ParquetWriter, Schema } from "parquetjs"; import { PassThrough } from "stream"; const s3Client = new S3Client({ region: "your-region" }); const BUCKET_NAME = "your-bucket-name"; const PREFIX = "parquet-data/"; export const handler = async (event: any) => { // 1. 定义Parquet Schema(必须与已有数据schema一致) const schema = new Schema({ id: { type: "INT32" }, name: { type: "UTF8" }, age: { type: "INT32" } }); // 2. 准备待写入的数据 const data = [ { id: 4, name: "David", age: 40 }, { id: 5, name: "Eve", age: 28 } ]; // 3. 创建Parquet写入流并转换为Buffer const passThrough = new PassThrough(); const writer = await ParquetWriter.openStream(schema, passThrough); for (const row of data) { await writer.appendRow(row); } await writer.close(); // 4. 生成唯一文件名,写入S3前缀下(实现逻辑追加) const fileName = `${PREFIX}data-${Date.now()}.parquet`; const command = new PutObjectCommand({ Bucket: BUCKET_NAME, Key: fileName, Body: passThrough.read() }); await s3Client.send(command); return { statusCode: 200, body: "数据成功写入/追加至S3 Parquet" }; };
单个Parquet文件追加的说明(不推荐)
若必须追加到单个Parquet文件,需先下载现有文件、读取数据、合并新数据后重新上传,此方式仅适用于小文件,示例逻辑如下:
- 使用
GetObjectCommand下载S3上的目标Parquet文件。 - 用
ParquetReader读取文件内容,合并新数据。 - 用
ParquetWriter生成新的Parquet文件,再用PutObjectCommand覆盖原文件。
由于Lambda内存和执行时间限制,此方式不适合大文件场景。
通用注意事项
- Lambda执行角色需配置S3的对应权限(读/写/列表)。
- 注意Parquet schema的一致性,追加时字段类型和数量需与已有数据匹配。
- 对于大规模数据,建议使用批量写入或调整Lambda的内存配置(更高内存对应更强的CPU和网络性能)。
内容的提问来源于stack exchange,提问作者rubenhak
相关产品推荐
相关产品推荐

