You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark Streaming 3.4.1 sourceArchiveDir无法归档文件问题咨询

Spark 3.4.1 FileStreamSource cleanSource归档路径重复问题临时解决办法

问题场景

使用Scala结合Spark 3.4.1进行流处理时,配置cleanSource为archive、sourceArchiveDir为archived,代码如下:

spark.readStream
  .option("cleanSource", "archive")
  .option("sourceArchiveDir", "archived")
  .option("enforceSchema", false)
  .option("header", includeHeader)
  .option("inferSchema", inferSchema)
  .options(otherOptions)
  .schema(csvSchema.orNull)
  .csv(FileUtils.getPath(sourceSettings.dataFolderPath, mappingSource.path).toString)

其中读取路径为相对路径test-data\streaming-folder\patients,但启动流后Spark未将源文件移动到归档目录。调试发现org.apache.spark.sql.execution.streaming.FileStreamSource.scala的cleanTask方法中,归档路径拼接错误:

  • 源文件路径:file:/C:/dev/be/data-integration-suite/test-data/streaming-folder/patients/patients-success.csv
  • 生成的错误归档路径:file:/C:/dev/be/data-integration-suite/archived/C:/dev/be/data-integration-suite/test-data/streaming-folder/patients/patients-success.csv
    根目录重复导致归档失败,正确路径应为file:/C:/dev/be/data-integration-suite/archived/test-data/streaming-folder/patients/patients-success.csv。

临时解决方案

1. 使用绝对URI格式的归档目录路径

将sourceArchiveDir配置为项目根目录下归档文件夹的绝对URI路径,避免Spark拼接路径时重复根目录:

spark.readStream
  .option("cleanSource", "archive")
  .option("sourceArchiveDir", "file:/C:/dev/be/data-integration-suite/archived/") // 绝对URI路径
  .option("enforceSchema", false)
  .option("header", includeHeader)
  .option("inferSchema", inferSchema)
  .options(otherOptions)
  .schema(csvSchema.orNull)
  .csv(FileUtils.getPath(sourceSettings.dataFolderPath, mappingSource.path).toString)

2. 自定义文件归档逻辑

关闭Spark自带的cleanSource归档,通过流监听或foreachBatch手动处理文件移动:

import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.spark.sql.streaming.StreamingQueryListener

// 跟踪已处理文件,避免重复操作
val processedFiles = scala.collection.mutable.HashSet[String]()

// 注册流监听获取已处理文件路径
spark.streams.addListener(new StreamingQueryListener() {
  override def onQueryStarted(event: StreamingQueryListener.QueryStartedEvent): Unit = {}
  override def onQueryProgress(event: StreamingQueryListener.QueryProgressEvent): Unit = {
    event.progress.sources.flatMap(_.endOffset.json.split("\"path\":\"([^\"]+)\""))
      .filter(_.startsWith("file:"))
      .foreach { filePath =>
        if (!processedFiles.contains(filePath)) {
          processedFiles.add(filePath)
          val fs = FileSystem.get(spark.sparkContext.hadoopConfiguration)
          val srcPath = new Path(filePath)
          // 构建目标路径:归档目录 + 源文件相对路径
          val relativePath = filePath.replace("file:/C:/dev/be/data-integration-suite/", "")
          val destPath = new Path("file:/C:/dev/be/data-integration-suite/archived/" + relativePath)
          // 创建目标目录(如果不存在)
          if (!fs.exists(destPath.getParent)) {
            fs.mkdirs(destPath.getParent)
          }
          // 移动文件
          if (fs.exists(srcPath)) {
            fs.rename(srcPath, destPath)
          }
        }
      }
  }
  override def onQueryTerminated(event: StreamingQueryListener.QueryTerminatedEvent): Unit = {}
})

// 构建流查询
spark.readStream
  .option("cleanSource", "off") // 关闭自带归档
  .option("enforceSchema", false)
  .option("header", includeHeader)
  .option("inferSchema", inferSchema)
  .options(otherOptions)
  .schema(csvSchema.orNull)
  .csv(FileUtils.getPath(sourceSettings.dataFolderPath, mappingSource.path).toString)
  .writeStream
  .foreachBatch { (df, batchId) =>
    // 执行你的业务处理逻辑
    df.write.mode("append").save("your-output-path")
  }
  .start()
  .awaitTermination()

3. 临时修改Spark源码编译

如果具备编译环境,可修改FileStreamSource.scala中cleanTask方法的路径拼接逻辑:

  • 找到构建newPath的代码,将源文件的绝对路径转换为相对于项目根目录的相对路径后,再与归档目录拼接。
  • 重新编译Spark的sql/core模块,替换项目依赖的对应jar包。

内容的提问来源于stack exchange,提问作者Y. Emre

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 16:33:43