PySpark中出现“Number of completed tasks exceeds total tasks”报错及大数据集资源问题求助
PySpark中出现“Number of completed tasks exceeds total tasks”报错及大数据集资源问题求助
兄弟我太懂你面对20亿行大数据集卡壳的崩溃了!直接硬刚全量数据就算取子集都能触发奇怪报错,你遇到的这个“Number of completed tasks exceeds total tasks”,还有资源吃紧的问题,多半是limit操作的隐藏坑加上默认配置顶不住大负载导致的,咱一步步来解决:
先唠唠为啥会出这个任务计数的报错
limit的隐形成本:你以为limit(10000000)是直接取前1000万行就完事?其实Spark的limit底层会先触发作业去扫描数据集的分区,要是原始数据分区极不均衡,或者集群资源不够导致任务反复重试,就会出现任务计数混乱的情况——重试的任务会被重复统计,最后就出现“完成任务数超过总任务数”的离谱报错。而且直接在20亿行的大DF上跑limit,本质还是要扫全量分区,资源消耗一点不少。- 默认配置顶不住大负载:用Jupyter跑PySpark的话,Driver和Executor的默认内存、核数都很低,面对大分区数据,Executor内存不够会频繁GC,任务反复失败重试,进一步加剧任务计数的异常。
给你几个亲测有用的解决步骤
1. 换个姿势抽取子集,别直接用limit硬刚
试试先采样或者重分区再取子集,能大幅减少资源消耗:
from pyspark.sql import functions as F # 方式1:先采样再limit,避免扫全量分区 # 20亿行的0.005就是1000万,seed固定保证结果可复现 df_sampled = df.sample(fraction=0.005, seed=42) df_10M = df_sampled.limit(10000000) # 方式2:先重分区减少小分区数量,让数据分布更均匀 # 分区数可以根据集群核数调整,比如100-200之间 df_repartitioned = df.repartition(100) df_10M = df_repartitioned.limit(10000000)
2. 调大Spark资源配置,别用默认的“乞丐版”
在Jupyter里初始化SparkSession的时候,手动加资源配置,比如:
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("BigDataSubsetDemo") \ .config("spark.driver.memory", "16g") # Driver内存根据你本地/集群节点情况调 .config("spark.executor.memory", "32g") # Executor内存别超过节点总内存的70% .config("spark.executor.cores", 8) # 每个Executor的核数,别抢满节点资源 .config("spark.executor.instances", 10) # 启动多少个Executor .config("spark.sql.shuffle.partitions", 100) # 默认2000太夸张,调小到和Executor核数匹配 .getOrCreate()
注意配置要贴合你的集群实际情况,别瞎开超资源,不然会被集群调度器干掉。
3. 优化后续的数据处理逻辑,少做无用功
你后面的withColumn堆一堆再drop旧列,不如直接用select精准选需要的列,能减少Spark加载的数据量:
# 别用withColumn加完再drop,直接select需要的列+新生成的列 df_a = df_10M.select( "col_you_need_1", "col_you_need_2", # 新列直接在select里定义 F.col("old_col").cast("int").alias("new_col1"), F.when(F.col("status") == "success", 1).otherwise(0).alias("new_col2") ) df_a.show(10)
这样Spark在扫描数据时只会加载你需要的列,内存压力能小很多。
4. 用Spark UI排查细节
打开Spark UI(默认是http://localhost:4040),去Jobs和Stages页面看看,是不是某个Stage的任务在反复失败重试?要是有失败任务,点进去看日志,是内存不够OOM,还是数据倾斜导致某个任务跑不动?数据倾斜的话可以给倾斜列加盐拆分,或者把大拆小处理。
你可以先从调整子集抽取方式和Spark配置入手,应该能先解决那个任务计数的报错,缓解资源压力。
内容来源于stack exchange
相关产品推荐
相关产品推荐

