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个分区。
- 解决措施:
- 显式重分区:在转换为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; }); - 调整分区大小阈值:读取CSV时通过配置减小分区大小,让Spark自动生成更多分区:
Dataset<Row> data = sparkSession.read() .option(HEADER, TRUE) .option("maxPartitionBytes", "1mb") .csv("path to csv"); - 验证分区数:添加代码确认分区数是否符合预期:
System.out.println("当前RDD分区数:" + oprdd.getNumPartitions());
- 显式重分区:在转换为JavaRDD后调用
内容的提问来源于stack exchange,提问作者Santhosh Kumar Sreeramoju
相关产品推荐
相关产品推荐

