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

向AWS Batch Transform Job高效流式传输Parquet数据的方案咨询

解决SageMaker Batch Transform流式处理大Parquet文件的方案

针对你遇到的Parquet文件无法像CSV那样用SplitType=Line拆分流式处理的问题,以下是两种无需手动拆分文件的可行方案:

方案1:利用Parquet行组(Row Group)配合SplitType='Record'

这是最接近CSV方案的简便方法,核心思路是在生成Parquet文件时控制行组大小≤10MB,让SageMaker Batch Transform自动按行组拆分流式传输。

步骤:

  1. 生成符合要求的Parquet文件
    在数据预处理阶段,使用PyArrow或Pandas生成Parquet时,设置row_group_size参数,让每个行组的大小接近但不超过10MB(可根据数据的平均行大小估算行组行数)。
    示例代码(PyArrow):

    import pyarrow as pa
    import pyarrow.parquet as pq
    import pandas as pd
    
    # 读取清洗后的数据
    df = pd.read_csv("cleaned_data.csv")
    table = pa.Table.from_pandas(df)
    
    # 估算行组大小:假设每行约1KB,10MB对应10000行(可根据实际数据调整)
    target_row_group_size = 10000
    pq.write_table(
        table,
        "preprocessed_inference.parquet",
        row_group_size=target_row_group_size,
        compression="none"  # 若使用压缩,需对应设置TransformInput中的CompressionType
    )
    
  2. 修改Batch Transform配置
    将ContentType改为Parquet的MIME类型,SplitType设为Record,SageMaker会自动按Parquet的行组拆分文件并流式传输:

    inference_job_config = {
        'TransformJobName': inference_job_name,
        'ModelName': self.model_name,
        'TransformInput': {
            'DataSource': {
                'S3DataSource': {
                    'S3Uri': f's3://{self.s3_bucket}/{self.s3_preprocessed_inference}/',
                    'S3DataType': 'S3Prefix'
                }
            },
            'ContentType': 'application/x-parquet',
            'CompressionType': 'None',  # 若用压缩则改为对应值(如'GZIP')
            'SplitType': 'Record'  # 按Parquet行组拆分
        },
        'TransformOutput': {
            'S3OutputPath': f's3://{self.s3_bucket}/{self.s3_inference_results}/',
            'AssembleWith': 'None',
            'Accept': 'application/x-parquet'  # 输出也设为Parquet格式
        },
        'TransformResources': {
            'InstanceType': self.instance_type,
            'InstanceCount': 1
        },
        'DataProcessing': {
            'JoinSource': 'None'
        }
    }
    

优点:

  • 无需自定义容器,仅需调整数据生成和Batch Transform配置,成本极低
  • 完全模拟CSV的SplitType=Line自动拆分逻辑,无需手动管理文件大小

方案2:自定义Batch Transform容器实现流式Parquet解析

如果不想限制Parquet文件的行组大小,可通过自定义Docker容器,让推理端点直接流式解析大Parquet文件,无需等待整个文件传输完成。

核心思路:

在自定义容器的推理代码中,使用PyArrow/FastParquet的流式读取能力,接收SageMaker发来的字节流后,逐行组读取并处理数据,边读边推理,避免内存溢出。

关键代码示例(推理函数):

import pyarrow.parquet as pq
import io
import pandas as pd

# 预先加载训练好的模型(根据你的模型类型调整)
model = load_your_model()

def predict(data, context):
    # 将HTTP接收的字节流转为可读取的对象
    parquet_stream = io.BytesIO(data)
    parquet_file = pq.ParquetFile(parquet_stream)
    
    results = []
    # 逐行组读取并处理
    for batch in parquet_file.iter_batches():
        df = batch.to_pandas()
        # 执行推理逻辑
        df["prediction"] = model.predict(df)
        results.append(df)
    
    # 合并结果并转为Parquet字节流返回
    final_df = pd.concat(results)
    output_buffer = io.BytesIO()
    final_df.to_parquet(output_buffer, index=False)
    return output_buffer.getvalue()

后续步骤:

  • 将推理代码、模型依赖打包成Docker镜像,推送到ECR
  • 在SageMaker中基于该镜像创建模型,再启动Batch Transform任务(无需设置SplitType,容器自行处理流式数据)

优点:

  • 完全不受Parquet文件大小限制,支持任意大文件的流式处理
  • 灵活性高,可自定义数据处理逻辑

缺点:

  • 需要构建和维护自定义Docker镜像,有一定开发和运维成本

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 05:34:58