Spark Structured Streaming AvailableNow Trigger故障时的结束偏移量行为咨询
Spark Structured Streaming AvailableNow Trigger 故障场景行为解析
初始结束偏移量是否会因故障变化?
不会。使用AvailableNow Trigger时,初始结束偏移量是查询启动瞬间Spark计算出的数据源当前最大可用偏移量——这是一个启动时的快照值,故障发生后这个初始值不会被修改。
故障恢复后的结束偏移量变化
启用检查点后,故障恢复时的结束偏移量遵循以下规则:
- 恢复时,Spark先从检查点读取上次故障前已成功处理完成的偏移量作为起始位置;
- 然后重新计算当前数据源的最新可用偏移量,将其作为本次恢复后的新结束偏移量;
- 简言之,恢复后的查询会处理「断点偏移量」到「恢复时数据源最新偏移量」之间的所有数据,包括故障期间新增的数据源内容。
举个实际场景:
- 查询启动时,数据源最大偏移量为100(初始结束偏移量),处理到偏移量50时发生故障;
- 故障期间数据源新增数据,偏移量涨到150;
- 恢复查询时,起始偏移量为50,新的结束偏移量为150,查询会处理50到150区间的数据。
是否区分故障与正常结束场景?
是的,Spark Structured Streaming会明确区分这两种场景:
- 正常结束:当查询完整处理完初始结束偏移量内的所有数据后,会在检查点中标记查询为「已完成」。再次启动查询时,会从上次完成的偏移量开始,重新计算当前数据源的最新可用偏移量作为新的结束偏移量,处理新的一批数据。
- 故障结束:故障发生时,检查点仅记录到已成功处理的偏移量,不会标记查询完成。恢复时查询不会沿用故障前的初始结束偏移量,而是基于断点重新适配当前数据源的最新状态,确保所有未处理和新增的数据都被覆盖。
内容的提问来源于stack exchange,提问作者MaatDeamon
相关产品推荐
相关产品推荐

