PySpark Streaming:memory输出流未检测到初始CSV数据问题求助
问题原因及解决办法
核心原因
Checkpoint的文件标记干扰
你给memorysink设置了checkpointLocation,如果这个目录之前运行过相同查询,Spark会在checkpoint元数据里标记源文件夹中的初始文件为「已处理」。重启查询时,Spark会跳过这些已标记文件,只响应新增数据。而你给consolesink没配置checkpoint,每次启动都会重新扫描整个目录,所以能立即读取初始文件。无效参数打乱初始化逻辑
你的memorysink查询里加了option('path', ...)——这个参数是给文件类sink(比如csv/parquet)用的,memorysink只负责将结果写入Spark内存表,完全不需要输出路径。这个多余参数可能干扰了Spark对sink类型的判断,导致初始文件扫描异常。
解决步骤
清理Checkpoint目录
删除C:/GP/music-listening-behaviour-streaming/data/01_raw/checkpoints下的所有内容,重新运行代码。这样Spark会从头开始扫描源文件夹,初始文件就能被检测到。移除多余的path参数
修改memorysink的代码,删掉无效的路径配置,调整后代码如下:
query_ = tscount.writeStream.format("memory")\ .outputMode("complete")\ .queryName("countusers")\ .option("checkpointLocation", "C:/GP/music-listening-behaviour-streaming/data/01_raw/checkpoints")\ .start()
- 主动验证内存表数据
启动查询后,执行spark.sql("SELECT * FROM countusers").show()查看内存表内容,不要仅依赖status或lastProgress输出——这些信息不一定会实时反映内存表的实际数据。
内容的提问来源于stack exchange,提问作者gaut
相关产品推荐
相关产品推荐

