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

AWS Wrangler分块读取JSON转Parquet:类型转换与文件合并问题

解决Lambda处理大JSON转Parquet的内存不足与分块问题

一、分块读取时的数据类型转换错误解决

当使用chunksize参数时,wr.s3.read_json返回的是生成器(generator),而非完整的DataFrame,因此无法直接对生成器调用astype()、访问columns属性。必须遍历生成器中的每个DataFrame分块,逐个处理类型转换:

全量转字符串的正确写法

for df_chunk in wr.s3.read_json(path=[s3_json_location], lines=True, orient='records', chunksize=1000):
    df_chunk = df_chunk.astype(str)
    # 添加后续处理逻辑(如写入Parquet)

指定列类型转换的正确写法

for df_chunk in wr.s3.read_json(path=[s3_json_location], lines=True, orient='records', chunksize=1000):
    # 转换column1为字符串类型
    if 'column1' in df_chunk.columns:
        df_chunk['column1'] = df_chunk['column1'].apply(str)
    # 填充column5空值并转换为整数类型
    if 'column5' in df_chunk.columns:
        df_chunk['column5'] = df_chunk['column5'].fillna(0).astype(int)
    # 添加后续处理逻辑

二、分块写入后合并为单个Parquet文件

分块写入会生成多个Parquet文件,需根据文件大小选择合适的合并方案:

方案1:内存允许时合并后写入

如果合并后的总数据量在Lambda内存限制范围内(例如大JSON压缩为Parquet后体积可控),可先将所有分块加载到内存合并,再写入单个文件:

import pandas as pd

# 收集所有处理后的分块
df_list = []
for df_chunk in wr.s3.read_json(path=[s3_json_location], lines=True, orient='records', chunksize=1000):
    df_chunk = df_chunk.astype(str)
    df_list.append(df_chunk)

# 合并所有分块为单个DataFrame
full_df = pd.concat(df_list, ignore_index=True)

# 写入单个Parquet文件(dataset=False确保生成单个文件)
wr.s3.to_parquet(df=full_df, path=single_parquet_path, dataset=False)

注意:此方法仅适用于合并后数据量不超过Lambda内存上限的场景,否则会再次触发内存不足。

方案2:使用AWS Glue/Athena合并(适合超大型文件)

对于Lambda内存无法承载的超大型文件,推荐用AWS Glue或Athena完成合并:

  • Glue作业:创建Glue ETL作业,读取分块生成的Parquet数据集,通过调整作业内存和并行度,将数据重新写入为单个Parquet文件。
  • Athena查询:创建外部表指向分块Parquet文件的存储路径,执行CREATE TABLE target_table WITH (format = 'PARQUET', location = 's3://target-path/') AS SELECT * FROM source_table;语句,将查询结果导出为单个Parquet文件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 01:13:30