如何关闭PySpark单任务探测作业,让mapPartitions任务并行启动?
PySpark关闭单任务探测作业的解决方案
问题描述
我有一个简单的PySpark mapPartitions作业,计划拆分至200个Worker并行运行,各任务相互独立,单个任务耗时2小时,因此希望所有任务从一开始就并行执行。但PySpark总会先运行一个单任务探测作业(Stage 1仅包含1个任务,仅处理200个分片的第0个分片),之后才部署全部200个Worker。已尝试设置spark.speculation=false、scheduler.mode=FAIR、增加Worker数量等方法,但该单任务探测行为仍存在。请问如何让PySpark直接运行纯映射作业,关闭该单任务探测行为?
已尝试的配置及代码
conf = SparkConf() conf.set("spark.executor.memory", "60g") conf.set("spark.yarn.executor.memoryOverhead", "20g") conf.set("spark.executor.memoryOverhead", "20g") conf.set("spark.executor.instances", "250") conf.set("spark.sql.shuffle.partitions", "200") conf.set("spark.default.parallelism", "200") conf.set("spark.dynamicAllocation.maxExecutors", "300") conf.set("executor-cores", '8') conf.set('spark.speculation', 'false') conf.set('spark.scheduler.mode', 'FAIR') spark = ( SparkSession.builder.enableHiveSupport() .config(conf=conf) .appName("My-App") .getOrCreate() ) num_workers = 200 data = [(t,) for t in range(num_workers)] schema = ['global_shard_idx'] input_df = spark.createDataFrame(data=data, schema=schema) input_df \ .repartition(num_workers) .rdd .mapPartitionsWithIndex(resolve_partition) .toDF(['partition_idx']) .write .format('orc') .option('header', 'true') .mode('overwrite') .saveAsTable('temp.dummy_table')
解决方案
禁用动态分区探测:这个单任务探测是Spark SQL写入表时的默认行为,用于扫描数据获取分区信息,添加以下配置关闭:
conf.set("spark.sql.sources.partitionColumnTypeInference.enabled", "false") conf.set("spark.sql.hive.convertMetastoreOrc", "false") # 针对ORC格式的专属配置简化DataFrame与RDD的转换:代码中先转RDD再转回DataFrame的操作可能触发额外探测,直接使用DataFrame的
mapPartitions替代RDD的mapPartitionsWithIndex:from pyspark.sql.types import IntegerType def process_partition(iterator): for row in iterator: yield (row.global_shard_idx,) input_df.repartition(num_workers) \ .mapPartitions(process_partition, schema=IntegerType()) \ .write.format('orc') \ .option('header', 'true') \ .mode('overwrite') \ .saveAsTable('temp.dummy_table')关闭统计信息采样:禁用Spark的初始采样行为,避免触发探测任务:
conf.set("spark.sql.statistics.histogram.enabled", "false") conf.set("spark.sql.statistics.fallBackToHadoop", "false")验证分区数量:确保
repartition(num_workers)确实生成200个分区,可在写入前执行print(input_df.repartition(num_workers).rdd.getNumPartitions())确认,避免数据分布异常导致分区合并。
内容的提问来源于stack exchange,提问作者John Jiang
相关产品推荐
相关产品推荐

