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
相关产品推荐
相关产品推荐

