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

Scala 2.13与Spark 3.2.1无法在Executor上并行执行问题求助

Spark本地模式单线程运行问题排查与解决

问题描述

使用Scala 2.13和Spark 3.2.1进行列表并行处理,已将Spark Master设置为local[4](4线程本地模式),但任务始终以单线程运行,所有输出均来自同一个任务线程,怀疑遗漏了关键配置。

相关代码与配置

打印方法

public static void m1(String i){
    System.out.println("===thread==="+Thread.currentThread().getName()+"===value==="+i);
}

Spark实现代码

final SparkSession sparkSession = SparkSession.builder()
                    .appName("appName")
                    .config("spark.sql.shuffle.partitions", 1500)
                    .config("spark.debug.maxToStringFields", 1500)
                    .config("spark.master", "local[4]")
                   .getOrCreate();
Dataset<Row> data = sparkSession.read().option(HEADER, TRUE).csv("path to csv");
Dataset<Row> lines = data.na().fill("");
  
JavaRDD<Row> oprdd = lines.javaRDD().map(x -> {
    try {
        m1(x.mkString());
    } catch (Exception e) {
        return  x;
    }
    return x;
});

oprdd.rdd().collect();

依赖配置

compile 'org.apache.spark:spark-core_2.13:3.2.1'
compile 'org.apache.spark:spark-sql_2.13:3.2.1'

输出结果

===thread==Executor task launch worker for task 0.0 in stage 2.0 (TID 5)===value==bcd
===thread==Executor task launch worker for task 0.0 in stage 2.0 (TID 5)===value==efg
===thread==Executor task launch worker for task 0.0 in stage 2.0 (TID 5)===value==hij
===thread==Executor task launch worker for task 0.0 in stage 2.0 (TID 5)===value==klm
===thread==Executor task launch worker for task 0.0 in stage 2.0 (TID 5)===value==nop

问题原因及解决办法

  • 核心原因:RDD分区数为1,Spark会为每个分区分配一个任务,单分区自然只会启动一个线程处理。CSV文件较小时,Spark默认按文件大小(默认128MB阈值)划分分区,小文件只会生成1个分区。
  • 解决措施:
    1. 显式重分区:在转换为JavaRDD后调用repartition(n)指定分区数(n对应线程数,比如4):
      JavaRDD<Row> oprdd = lines.javaRDD().repartition(4).map(x -> {
          try {
              m1(x.mkString());
          } catch (Exception e) {
              return  x;
          }
          return x;
      });
      
    2. 调整分区大小阈值:读取CSV时通过配置减小分区大小,让Spark自动生成更多分区:
      Dataset<Row> data = sparkSession.read()
          .option(HEADER, TRUE)
          .option("maxPartitionBytes", "1mb")
          .csv("path to csv");
      
    3. 验证分区数:添加代码确认分区数是否符合预期:
      System.out.println("当前RDD分区数:" + oprdd.getNumPartitions());
      

内容的提问来源于stack exchange,提问作者Santhosh Kumar Sreeramoju

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 09:15:35