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

Azure Databricks环境下Spark任务无法正常并行执行问题咨询

问题根因

你的任务并行度由groupby(['Material'])生成的task数量决定,Azure Databricks默认的小数据集分区策略、自动扩缩容触发逻辑和自建Spark存在差异,会导致初始分区数极低、集群扩容不及时,最终所有计算串行执行。

针对性配置修改

1. Spark核心并行度参数配置

创建SparkSession时新增以下配置,匹配你的集群资源和任务量:

spark = SparkSession \
    .builder \
    .appName("test") \
    .config("spark.sql.execution.arrow.pyspark.enabled", "true") \
    # 新增以下两行,shuffle分区数匹配你的SKU总量(约6000),保证每个SKU对应独立task
    .config("spark.sql.shuffle.partitions", "6000") \
    # 默认并行度匹配集群总核心数的2倍:25worker*4核*2=200
    .config("spark.default.parallelism", "200") \
    .getOrCreate()

2. 手动重分区保证任务打散

创建Spark DataFrame后、groupby操作前,按分组键重分区,避免所有数据集中在少数分区:

df_spark = spark.createDataFrame(df).repartition(6000, "Material")

3. 调整集群动态资源分配参数

在Azure Databricks集群配置页的「Spark配置」栏添加以下参数,加快集群扩容触发速度:

spark.dynamicAllocation.enabled true
spark.dynamicAllocation.minExecutors 2
spark.dynamicAllocation.maxExecutors 25
spark.dynamicAllocation.schedulerBacklogTimeout 1s

代码优化建议

  • 修正代码笔误:原代码中df = df_OTB[['Material', 'Alpha']]的df_OTB未定义,应该修改为df = df[['Material', 'Alpha']]
  • 取消不必要的driver端数据拉取:原代码先collect再转Pandas写CSV的逻辑会把全量数据拉到driver节点,直接用Spark原生写CSV接口性能更高:
df_spark \
    .groupby(['Material']) \
    .applyInPandas(main, schema=schema) \
    .write.mode("overwrite") \
    .option("header", "true") \
    .csv("/mnt/simulaciones/Resultados/test.csv")
  • 全局变量DATA需要广播,否则worker节点执行时无法读取路径,建议将路径也封装为广播变量传递。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 15:27:00