Spark Autoloader检查点会保存哪些文件?基于过滤逻辑的技术问询
Spark流处理检查点与文件过滤问题解析
检查点实际保存的内容
你的代码里,检查点会存所有被Spark扫到的文件路径和它们的元数据(比如修改时间、文件大小),不是只存过滤后的文件。原因很简单:
- CloudFiles这个数据源的文件追踪是在
load()读取文件的时候做的,这时候它会把指定路径下所有符合扫描规则的文件(因为开了recursiveFileLookup,所以包括子文件夹)都记下来,不管你后面有没有用filter把旧文件筛掉。 filter是读完文件之后才执行的过滤操作,检查点根本不知道这个步骤,它只认读取阶段扫到的所有文件。
带来的问题就是:
- 第一次跑的时候,Spark会扫全量文件,过滤掉旧的再处理新的,但检查点会把所有扫过的文件都存起来。
- 之后再启动任务,Spark会跳过检查点里记录的所有文件——哪怕某个旧文件后来被修改到符合时间条件了,也不会再被读取处理。
满足你需求的正确做法
你得在读取文件的源头就过滤掉旧文件,而不是读完再筛。CloudFiles本身就支持通过参数直接按修改时间过滤,这样检查点里只会存符合条件的文件,同时避免每次都扫全量旧文件。
修改后的代码如下:
val streamingQuery = spark.readStream .format("cloudFiles") .option("cloudFiles.format", "binaryFile") .schema("`path` STRING, `modificationTime` TIMESTAMP, `length` BIGINT, `content` BINARY") .option("recursiveFileLookup", "true") // 直接在读取源头过滤:只取修改时间晚于指定时间的文件 .option("cloudFiles.modificationTimeRangeStart", "2023-10-30 07:00:00") .load("my path here") .writeStream .trigger(Trigger.AvailableNow()) .foreachBatch (my code goes here on how I copy files) .option("checkpointLocation", "my path here") .start() .awaitTermination()
这么做的好处:
- 不用扫全量旧文件:读取阶段直接跳过修改时间早于指定值的文件,不会把这些旧文件的信息写到检查点里。
- 检查点只记需要处理的文件:下次启动任务时,Spark只会扫描新出现的、符合时间条件的文件,完全不用碰旧文件。
- 旧文件修改后能被重新处理:如果某个旧文件后来被修改,且修改后的时间符合要求,CloudFiles会检测到文件元数据变化,重新读取处理它。
额外提示
要是需要指定时间范围的结束点,还可以加cloudFiles.modificationTimeRangeEnd参数。另外要注意,这些参数需要Spark 3.1以上版本搭配Delta Lake 2.0以上的CloudFiles才能支持。
内容的提问来源于stack exchange,提问作者Tamás Godányi
相关产品推荐
相关产品推荐

