Java Spark中如何在Executor端实现数据写入S3 Parquet文件?
这个问题我之前做Spark任务时也踩过坑!核心原因是没搞清楚Spark Driver和Executor的职责边界,先给你拆解下报错的根源,再给你几个从优到次的解决方案:
为什么会抛出parentSessionState() is null异常?
SparkSession是Driver端专属的上下文对象,它的生命周期完全由Driver节点管理。Spark的架构设计里,Executor节点只负责执行Driver分发的具体任务,不维护独立的Spark上下文——你在foreachPartition里调用dummyClass.createParquet()时,用到的SparkSession实例被传到了Executor环境,但这个实例在Executor上是无效的,没有对应的父会话状态,所以才会抛出这个异常。而且Executor节点也不允许创建新的SparkSession,这是Spark的架构限制。
解决方案:从最优到次优选择
方案一:用Spark原生分布式写入(最推荐,符合Spark设计理念)
Spark的Dataset/DataFrame.write是天生的分布式操作,会自动把数据分发到各个Executor并行写入S3,完全不需要你在Executor里手动处理Session。这不仅性能最优,还能避免大量小文件的问题,是Spark任务的标准写法。
调整思路:把业务逻辑移到Dataset层面,统一写入
如果你的业务逻辑可以在Dataset上通过Spark的API实现,直接改造Processing类的run方法即可:
public void run() { // 保留原来的fetching data和数据准备逻辑 Dataset<Row> myData = ...; // 已经处理好的原始数据集 // 用Spark的Dataset API完成你原来在每个Row里的业务逻辑 // 比如生成字段、转换结构,替代原来的row.getAs和dummyClass的处理 Dataset<Row> finalData = myData .withColumn("newCol", functions.col("colName")) // 这里替换成你的业务逻辑 .select("col1", "col1w"); // 匹配你要写入的Schema // 直接用Spark原生write写入S3,自动分布式执行 String s3Path = "s3a://bucket/key/key/"; finalData.write() .mode("overwrite") .format("parquet") .save(s3Path); }
如果需要复杂的单Row/单分区处理:用mapPartitions转换后统一写入
如果你的业务逻辑必须调用外部服务、做复杂计算(不能用Spark API实现),可以用mapPartitions在每个Executor的分区内批量处理,再统一写入:
public void run() { Dataset<Row> myData = ...; // 用mapPartitions处理整个分区的数据(每个分区在Executor上执行一次) Dataset<Row> processedData = myData.mapPartitions(partition -> { // 每个分区初始化一次DummyClass,避免重复创建实例 DummyClass dummy = new DummyClass(...); List<Row> resultRows = new ArrayList<>(); while (partition.hasNext()) { Row row = partition.next(); String myField = row.getAs("colName"); // 调用DummyClass的方法处理数据,返回要写入的Row Row processedRow = dummy.processRow(myField); // 你需要新增这个方法,返回Row而非直接写入 resultRows.add(processedRow); } return resultRows.iterator(); }, createFinalSchema()); // 传入最终要写入的Schema // 统一写入S3 processedData.write() .mode("overwrite") .parquet("s3a://bucket/key/key/"); } // 定义最终的Parquet Schema private StructType createFinalSchema() { return new StructType() .add("col1", DataTypes.StringType, false) .add("col1w", DataTypes.StringType, false); }
这种方式既保留了你的业务逻辑,又符合Spark的分布式设计,性能和稳定性都有保障。
方案二:用Hadoop ParquetWriter直接写入(仅适用于特殊业务场景)
如果你的业务必须要给每个小批次数据单独写入不同的S3路径(比如按myField分目录),那可以绕过SparkSession,直接用Hadoop的ParquetWriter在Executor上写入——这个方式不需要依赖Spark上下文,直接操作文件系统。
调整DummyClass的createParquet方法:
public void createParquet(String myField) { List<Row> rowVals = new ArrayList<>(); StructType schema = createSchema(); // 保留原来的populate rowVals逻辑 String s3Path = "s3a://bucket/key/key/" + myField + "/"; // 配置S3和Parquet相关参数 Configuration conf = new Configuration(); conf.set("fs.s3a.endpoint", "s3-us-east-1.amazonaws.com"); // 按需添加其他S3配置,比如access key、secret等 conf.set("fs.s3a.access.key", "your-access-key"); conf.set("fs.s3a.secret.key", "your-secret-key"); // 把Spark StructType转换成Parquet需要的MessageType MessageType parquetSchema = new SparkToParquetSchemaConverter().convert(schema); // 初始化ParquetWriter Path path = new Path(s3Path); ParquetWriter<GenericRecord> writer = AvroParquetWriter.<GenericRecord>builder(path) .withSchema(parquetSchema) .withConf(conf) .withWriteMode(ParquetFileWriter.Mode.OVERWRITE) .build(); // 把Spark Row转换成Avro GenericRecord并写入 for (Row row : rowVals) { GenericRecord record = new GenericData.Record(parquetSchema); record.put("col1", row.getAs("col1")); record.put("col1w", row.getAs("col1w")); writer.write(record); } // 务必关闭writer,避免资源泄漏 writer.close(); }
⚠️ 注意:这种方式会在Executor上生成大量小文件,严重影响S3的存储效率和后续查询性能,后续必须做小文件合并(比如用Spark的coalesce或repartition,或者S3的批量操作),所以除非业务必须,否则强烈不推荐。
最后总结
永远不要在Executor里尝试使用SparkSession进行写入操作,这完全违反了Spark的架构设计。优先选择方案一,利用Spark的原生分布式写入能力,既高效又省心;只有在特殊业务场景下,再考虑方案二,同时做好小文件治理。
备注:内容来源于stack exchange,提问作者Joe

