处理24GB S3数据集的PySpark集群理想配置方案
问题背景
处理目标数据集路径为 s3://commoncrawl/crawl-001/2008/06/19/1/,总大小24GB,需求为筛选出所有text/html类型的请求,将结果保存到自有S3存储桶。
已遇到的报错
- 内存溢出报错:
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小时左右抛出
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%左右的内存占用
- 新增executor内存overhead配置:
- 解决会话超时问题:
Session isn't active是会话心跳超时被回收导致,新增两个配置即可解决:spark.network.timeout=36000spark.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
相关产品推荐
相关产品推荐

