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

读取S3中999个GZ压缩JSON文件并合并为Pandas DataFrame求助

解决方案

核心问题修复与完整代码

import pandas as pd
from pyspark.sql.functions import lower, col, length

# 初始化空DataFrame用于存放最终合并结果
final_pandas_df = pd.DataFrame()

# 遍历0-999生成所有文件名
for i in range(1000):
    # 生成带5位前导零的文件名,适配part-00000.gz~part-00999.gz格式
    file_path = f"s3a://my_bucket/company_v20_dl/part-{i:05d}.gz"
    
    # 读取并处理数据(Spark层先完成所有转换,减少Pandas端数据量)
    spark_df = spark.read.json(file_path)
    processed_spark_df = spark_df.select('id', 'summary', 'website') \
                                 .withColumn('text', lower(col('summary'))) \
                                 .select('id', 'text', 'website') \
                                 .withColumn("text_length", length("text"))
    
    # 转为Pandas DataFrame
    pandas_df = processed_spark_df.toPandas()
    
    # 打印当前文件处理结果(脚本环境必须用print才能显示)
    print(f"===== 处理文件part-{i:05d}.gz =====")
    print(pandas_df.head())
    
    # 合并到最终大表
    final_pandas_df = pd.concat([final_pandas_df, pandas_df], ignore_index=True)

# 查看合并后的完整数据
print("\n===== 合并后的完整数据 =====")
print(final_pandas_df.head())

关键问题说明

  1. 文件名自动补零
    使用Python格式化语法i:05d,将整数i格式化为5位固定长度的字符串,自动补前导零。比如i=1会生成00001,i=100生成00100,完美匹配你的文件名规则。

  2. Pandas DataFrame显示与合并

    • 脚本环境中pandas_df.head()不会自动输出,必须用print()才能看到结果;如果是Jupyter Notebook,直接写pandas_df.head()即可显示。
    • 通过pd.concat逐次合并每个文件的Pandas DataFrame,ignore_index=True避免合并后出现重复索引。

额外优化建议

  • 如果单个文件转Pandas仍存在内存压力,可在Spark层先做数据过滤(比如剔除summary为空的行),进一步减少转入Pandas的数据量。
  • 批量合并:可以每处理100个文件就合并一次到临时表,再与最终表合并,减少pd.concat的调用次数,提升效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 06:45:37