如何获取Spark append模式下DataFrame写入操作生成的文件列表
解决方案
针对EMR环境Spark append模式写入S3后获取单次写入新增文件的需求,以下是3种可落地的实现方式:
方案1:使用Spark QueryExecutionListener 监听写入事件(最推荐)
Spark内置的查询执行监听器可以在写入任务完成后直接返回本次写入生成的所有文件路径,无需扫描S3,无额外开销,结果100%准确。
实现代码:
import org.apache.spark.sql.execution.QueryExecution import org.apache.spark.sql.util.QueryExecutionListener import scala.collection.mutable.ListBuffer // 定义监听器 val writeFileListener = new QueryExecutionListener { override def onSuccess(funcName: String, qe: QueryExecution, durationNs: Long): Unit = { // 提取写入的文件路径 val outputFiles = qe.executedPlan.collectLeaves().flatMap { case write: org.apache.spark.sql.execution.datasources.WriteFilesExec => write.writeFilesSpec.map(_.outputFiles).getOrElse(Seq.empty) case _ => Seq.empty } // 这里可以把文件列表存到你需要的地方,比如变量、数据库等 if (outputFiles.nonEmpty) { println(s"本次写入生成的文件列表:${outputFiles.mkString(",")}") } } override def onFailure(funcName: String, qe: QueryExecution, exception: Exception): Unit = {} } // 注册监听器到SparkSession spark.listenerManager.register(writeFileListener) // 执行你的写入逻辑即可,监听器会自动捕获写入的文件 spark.read.json(Seq(jsonString).toDS) .write .mode("append") .json("s3://bucket/test")
方案2:自定义写入文件前缀,通过前缀区分
给每次写入任务配置唯一的文件前缀,后续直接通过S3前缀过滤就能拿到本次写入的所有文件,适合需要长期追溯写入任务的场景。
实现代码:
import org.apache.spark.sql.functions.lit // 生成本次写入的唯一标识,也可以用业务主键比如任务ID、时间戳等 val writeUuid = java.util.UUID.randomUUID().toString.replace("-", "") // 配置写入文件名前缀 spark.conf.set("mapreduce.output.basename", s"write-${writeUuid}-") // 执行写入 spark.read.json(Seq(jsonString).toDS) .withColumn("modified", lit("yes")) .write .mode("append") .json("s3://bucket/test") // 后续需要拿本次文件时,直接过滤S3路径下前缀为 write-${writeUuid}- 的文件即可
注意:EMR环境默认支持该参数,无需额外配置
方案3:写入前后扫描S3路径对比差异
如果不想修改现有写入逻辑,可在写入操作前后分别扫描目标路径的文件列表,取差集即为本次新增文件,适合小数据量目录的临时场景。
实现代码:
import org.apache.hadoop.fs.{Path, FileSystem} import java.net.URI val s3Path = new Path("s3://bucket/test") val fs = FileSystem.get(URI.create("s3://bucket/test"), spark.sparkContext.hadoopConfiguration) // 写入前的文件列表 val beforeFiles = fs.listFiles(s3Path, true).toSeq.map(_.getPath.toString).toSet // 执行写入逻辑 spark.read.json(Seq(jsonString).toDS) .write .mode("append") .json("s3://bucket/test") // 写入后的文件列表 val afterFiles = fs.listFiles(s3Path, true).toSeq.map(_.getPath.toString).toSet // 差集即为本次新增文件 val newFiles = afterFiles -- beforeFiles println(s"本次写入新增文件:${newFiles.mkString(",")}")
注意:如果目标路径有其他并发写入任务,该方法结果会不准确
内容的提问来源于stack exchange,提问作者zachd1_618
相关产品推荐
相关产品推荐

