Tomcat9运行Spark Driver:Java并行流加载Parquet文件异常求助
问题分析与解决方案
这个问题的核心原因很明确:SparkSession(以及底层的SparkContext)不是线程安全的,你用Java parallelStream 在多线程环境下并发创建/获取SparkSession,会触发类加载冲突,导致Parquet相关的内部类无法被正确加载——这就是你看到ClassNotFoundException的根源。而顺序流是单线程执行,所有操作都在同一个线程上下文里,类加载环境一致,所以不会出问题。
下面给出两种解决方案,优先推荐第一种(符合Spark的设计理念):
方案一:利用Spark自身的并行能力(推荐)
Spark本身就是为分布式并行计算设计的,完全不需要用Java的parallelStream来实现并行。我们可以复用单个SparkSession,一次性加载所有文件,然后通过Spark的内置函数统计每个文件的记录数,这是最高效且安全的方式。
代码示例:
// 只创建一次SparkSession,全局复用 SparkSession spark = get_spark_session(); List<String> filePaths = Arrays.asList( "hdfs://NNcluster/finalsnappy/2019/4/8291/table_8291_69_2019_04_01", "hdfs://NNcluster/finalsnappy/2019/4/8291/table_8291_69_2019_04_02" ); // 读取所有Parquet文件,并添加一列记录当前行所属的文件路径 Dataset<Row> fullDs = spark.read().format("parquet") .load(filePaths.toArray(new String[0])) .withColumn("file_path", functions.input_file_name()); // 按文件路径分组,统计每个文件的记录数 fullDs.groupBy("file_path") .count() .show();
这种方式的优势:
- 避免了多线程操作SparkSession的线程安全问题
- 利用Spark的分布式并行能力,比Java本地并行更高效
- 减少了Session创建的资源开销(SparkSession是重量级对象)
方案二:线程本地存储隔离SparkSession(不推荐)
如果你因为某些特殊原因必须使用Java的parallelStream,可以用ThreadLocal为每个线程创建独立的SparkSession,确保每个线程的Session完全隔离。但注意:SparkSession是重量级对象,多线程创建会消耗大量资源,不建议在生产环境使用。
代码示例:
// 用ThreadLocal为每个线程创建独立的SparkSession ThreadLocal<SparkSession> threadLocalSpark = ThreadLocal.withInitial(() -> get_spark_session()); List<String> list = Arrays.asList( "hdfs://NNcluster/finalsnappy/2019/4/8291/table_8291_69_2019_04_01", "hdfs://NNcluster/finalsnappy/2019/4/8291/table_8291_69_2019_04_02" ); list.parallelStream().forEach(fileName -> { SparkSession spark = threadLocalSpark.get(); Dataset<Row> tmpDS = spark.read().format("parquet").load(fileName); tmpDS.show(); System.out.println(String.format("File %s record count: %d", fileName, tmpDS.count())); }); // 最后清理线程本地存储的Session,避免资源泄漏 threadLocalSpark.remove();
关键注意事项
- SparkSession和SparkContext都是单例设计且线程不安全,永远不要在多线程环境下并发创建或访问它们
- 尽量复用同一个SparkSession完成所有任务,这是Spark的最佳实践
- 如果你之前的
get_spark_session()方法是每次都创建新的Session,那在并行流里会生成多个Session,这不仅会导致类加载问题,还会严重浪费集群资源
内容的提问来源于stack exchange,提问作者Yella Goud
相关产品推荐
相关产品推荐

