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
相关产品推荐
相关产品推荐

