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

如何关闭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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 11:43:13