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

处理24GB S3数据集的PySpark集群理想配置方案

问题背景

处理目标数据集路径为 s3://commoncrawl/crawl-001/2008/06/19/1/,总大小24GB,需求为筛选出所有text/html类型的请求,将结果保存到自有S3存储桶。

已遇到的报错
  1. 内存溢出报错:

Reason: Container killed by YARN for exceeding memory limits. 11.1 GB of 11 GB physical memory used.
初始集群配置为1主节点+2台m5.xlarge从节点,将全节点升级为m5.2xlarge后内存报错仍然存在。

  1. 扩容后新增报错:任务运行1小时左右抛出Session isn't active错误。
现有核心代码
# 读取数据逻辑
rdd=sc.wholeTextFiles('s3://commoncrawl/crawl-001/2008/06/19/1/')
# 结果保存逻辑
def toCSVLine(data):
  return ','.join(str(d) for d in data)

results = finalRdd.map(toCSVLine)
results.saveAsTextFile(
    path="s3://mybucket/folder/results/pages.csv",
    compressionCodecClass="org.apache.hadoop.io.compress.GzipCodec"
)
解决方案

1. 优先优化代码与配置(不需要盲目提升集群规格)

你遇到的内存问题核心不是集群规格不够,是wholeTextFiles的使用不符合Common Crawl数据的读取逻辑:wholeTextFiles会把单个文件的完整内容作为单条RDD记录加载到内存,Common Crawl的单分片文件体积很大,会直接导致单个executor内存被打满。

  • 替换读取逻辑:改用Common Crawl专用的WARC/ARC读取工具处理原始文件,或使用sc.binaryFiles配合流式解析逻辑读取,避免一次性加载全文件内容到内存
  • 调整RDD分区:读取数据后立即执行repartition(60),将全量数据拆分为60个左右的小分区,保证单分区数据量不超过200MB,降低单executor的内存压力
  • 调整Spark内存配置:
    • 新增executor内存overhead配置:spark.executor.memoryOverhead=4096(按executor内存的20%~30%设置即可)
    • 启用Kryo序列化:spark.serializer=org.apache.spark.serializer.KryoSerializer,可降低30%左右的内存占用
  • 解决会话超时问题:Session isn't active是会话心跳超时被回收导致,新增两个配置即可解决:
    • spark.network.timeout=36000
    • spark.executor.heartbeatInterval=30000
      若使用托管Spark服务(如EMR、Databricks)还要同步调整对应会话管理组件(如Livy)的超时阈值。

2. 合理集群规格

做完上述优化后,1主节点+3台m5.xlarge从节点即可稳定跑完24GB规模的Common Crawl筛选任务,不需要更高规格的节点。如果希望提升运行速度,可以最多扩展到5台m5.xlarge从节点,更高的配置只会造成资源浪费。

3. 写入优化

保存结果前可执行coalesce(10)将最终结果合并为10个左右的压缩分片,避免产生大量小文件,同时降低写入S3的请求开销。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 22:06:03