Spark availableNow模式下,流处理期间新增数据是否会被处理?
Spark Structured Streaming availableNow模式疑问解答
核心问题解答
availableNow模式的核心逻辑是一次性处理启动时刻所有可用的数据:启动时会对指定路径下的现有数据做快照,之后以微批方式逐步处理这些快照数据,处理完成后流任务会自动终止。
在处理现有数据的过程中,新上传的文件不在启动时的快照范围内,因此不会被本次运行的流任务处理。如果要处理这些新数据,需要重新启动一个使用availableNow模式的流任务。
你的代码问题说明
你连续调用了两次.trigger()配置,Spark的流任务只会生效最后一次的trigger设置,也就是代码里的processingTime='5 seconds'(固定间隔微批模式),前面的availableNow=True会被完全覆盖,这和你的预期行为不符。
如果要启用availableNow模式,只需保留一行trigger配置即可,修正后的代码如下:
bronze_query = (spark.readStream .format("cloudFiles") .option("cloudFiles.format", "json") .schema(schema) .load(<some_path>) .writeStream .format("delta") .outputMode("append") .trigger(availableNow=True) .option("checkpointLocation", f"<path>") .table("bronze"))
内容的提问来源于stack exchange,提问作者Elm662
相关产品推荐
相关产品推荐

