本地Java Spark程序并行异常及parallelize方法报错求助
问题1:程序仅单线程运行的解决方法
你设置的master("local[*]")本身是正确的,但程序单线程运行大概率是输入文件的分区数不足导致的:
- Spark的
textFile方法默认分区数是min(spark.default.parallelism, 2),如果你的文件很小(比如只有几行),Spark会只创建1个分区,自然只会用一个线程执行任务。 - 解决办法:读取文件时手动指定分区数,比如:
// 直接在textFile第二个参数指定分区数,比如和你的CPU线程数一致 JavaRDD<String> inputData = sc.textFile("src/main/resources/names.txt", 12); - 另外可以通过以下代码验证当前默认并行度:
如果输出不是12,可以在SparkSession/SparkConf中显式设置:System.out.println(spark.sparkContext().defaultParallelism());// SparkSession方式 SparkSession spark = SparkSession .builder() .master("local[*]") .appName("JavaWordCounter.com") .config("spark.default.parallelism", 12) .getOrCreate(); // 或者SparkConf方式 SparkConf conf = new SparkConf().setAppName("newSpark") .setMaster("local[*]") .set("spark.default.parallelism", "12");
问题2:parallelize处理文件RDD报错的解决方法
这个错误是因为你完全用错了parallelize方法:
parallelize的作用是将本地Java集合(比如List、Array)转换为RDD,而sc.textFile("...")已经直接返回了JavaRDD<String>,不需要再用parallelize处理。- 如果你想把文件RDD的内容转换成
JavaRDD<Integer>(比如每个字符串的长度),直接用map算子即可:SparkConf conf = new SparkConf().setAppName("newSpark").setMaster("local[*]"); JavaSparkContext sc = new JavaSparkContext(conf); JavaRDD<String> inputData = sc.textFile("src/main/resources/names.txt"); // 示例:将每个字符串转为其长度的Integer JavaRDD<Integer> myRdd = inputData.map(String::length); - 你之前用List调用
parallelize能运行,是因为List符合方法要求的参数类型,而RDD不符合,所以报错。
内容的提问来源于stack exchange,提问作者Arman
相关产品推荐
相关产品推荐

