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

在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文件,需先下载现有文件、读取数据、合并新数据后重新上传,此方式仅适用于小文件,示例逻辑如下:

  1. 使用GetObjectCommand下载S3上的目标Parquet文件。
  2. 用ParquetReader读取文件内容,合并新数据。
  3. 用ParquetWriter生成新的Parquet文件,再用PutObjectCommand覆盖原文件。
    由于Lambda内存和执行时间限制,此方式不适合大文件场景。

通用注意事项

  • Lambda执行角色需配置S3的对应权限(读/写/列表)。
  • 注意Parquet schema的一致性,追加时字段类型和数量需与已有数据匹配。
  • 对于大规模数据,建议使用批量写入或调整Lambda的内存配置(更高内存对应更强的CPU和网络性能)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 04:15:34