Spark Streaming任务处理HDFS文件异常:文件未消费且新文件无法识别
解决Spark监控本地Windows目录处理JSON文件时不消费、不识别新文件的问题
看起来你在测试Spark处理HDFS JSON文件时,先用本地Windows目录模拟,但碰到了文件没被消费、新增文件无法识别的问题——我之前做类似测试时也踩过不少坑,给你几个针对性的排查和解决思路:
一、先补全核心代码逻辑(你的代码片段不完整)
如果是用Spark Structured Streaming(现在推荐的方式),完整的本地文件监控处理代码应该是这样的,注意几个关键配置:
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.SparkSession; import org.apache.spark.sql.streaming.StreamingQuery; public class SparkConsumer { public static void main(String[] args) throws InterruptedException { execute(); } public static void execute() throws InterruptedException { // 本地测试必须用多线程,单线程会阻塞监控 SparkSession spark = SparkSession.builder() .master("local[*]") .appName("JsonFileConsumer") .getOrCreate(); // 读取JSON文件流,注意Windows路径格式 Dataset<Row> jsonDF = spark.readStream() .format("json") // Windows路径要用file:///前缀,双反斜杠或正斜杠 .option("path", "file:///C:/test/json_files") // 每次触发处理1个文件,方便测试观察 .option("maxFilesPerTrigger", 1) // 只处理.json后缀的文件,过滤临时文件 .option("pathGlobFilter", "*.json") // 处理完的文件自动删除,确认是否被消费 .option("cleanSource", "delete") .load(); // 必须输出到sink(控制台/存储),否则流不会启动 StreamingQuery query = jsonDF.writeStream() .outputMode("append") .format("console") .start(); // 阻塞程序保持运行,等待新文件 query.awaitTermination(); } }
二、Windows系统特有的坑点排查
路径格式错误
Windows的本地路径必须加上file:///前缀,路径分隔符用双反斜杠\\或者正斜杠/,比如file:///C:\\test\\json_dir,如果直接写C:\test\json_dir,Spark会识别成HDFS路径,自然找不到文件。文件未完全写入
Windows的文件锁机制会导致:如果JSON文件还在被其他程序写入(比如还在生成中,没有关闭输出流),Spark会判定文件不完整,不会消费。测试时可以:- 确保文件是完全生成好的(比如手动复制已完成的JSON文件到目录)
- 用
option("fileSuffix", ".done"),只监控带.done后缀的文件——生成完JSON后重命名为xxx.json.done,Spark才会处理。
单线程模式阻塞监控
本地测试时master必须设为local[*]或者local[2]以上,单线程local[1]会让Spark同时只能做一件事,监控线程被处理线程阻塞,根本无法识别新文件。权限与目录问题
- 确保Spark运行的用户(比如IDE的运行用户)有该目录的读写权限,避免用系统保护目录(如
C:\Windows下的文件夹) - 过滤临时文件:Windows生成文件时会有
.tmp这类临时文件,用pathGlobFilter只监控.json文件,避免干扰。
- 确保Spark运行的用户(比如IDE的运行用户)有该目录的读写权限,避免用系统保护目录(如
三、确认流是否正常运行
- 必须调用
query.awaitTermination(),否则程序启动流后会立刻退出,看起来像没消费文件 - 查看IDE控制台的Spark日志,搜索
FileStreamSource关键词,如果看到Found new files: [xxx.json]的日志,说明监控正常;如果没有,说明路径或配置有问题。
内容的提问来源于stack exchange,提问作者wandermonk
相关产品推荐
相关产品推荐

