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
相关产品推荐
相关产品推荐

