【发布时间】:2019-12-19 17:07:10
【问题描述】:
如果 Flink 应用程序在失败后启动备份或更新,是否保留不明确属于 KeyedState 或 OperatorState 的类变量?
例如,Flink 文档中描述的 BoundedOutOfOrdernessGenerator 有一个 currentMaxTimestamp 变量。如果 Flink 应用程序更新,currentMaxTimestamp 中的值会丢失,还是会写入应用程序更新之前创建的保存点?
这样做的真正原因是我想实现一个自定义水印生成器(similar to this),如果源空闲时间过长,它会在生成水印时切换到处理时间。但是,我希望根据类变量重置为其原始默认值(例如我上面提供的链接的示例中的 Long.MIN_VALUE),检测到应用程序在更新或失败后重新上线。这样,我可以确保水印生成器不会将耗时 5 分钟的应用程序更新误认为源空闲 5 分钟。
此外,如果应用程序更新,Flink 是否会重新启动每个水印生成器算子,即使水印生成器没有发生任何更改?
【问题讨论】:
标签: apache-flink flink-streaming