【问题标题】:How does Apache Beam Fault Tolerance work for Global Windows?Apache Beam 容错如何适用于全局 Windows?
【发布时间】:2019-05-21 01:13:09
【问题描述】:

我正在使用 Beam Python 构建管道。我有来自 PubSub 的带有 userId 和 buttonId 的事件流。我有一个全局窗口,用于维护所有用户单击按钮的次数。

如果一段时间后服务器重新启动运行 Direct Runner/Flink Runner,全局 windows 状态是否会恢复到管道中?

Beam 中的容错是如何工作的?

如何跟踪 PubSub 的偏移量/检查点?

Beam documentation 声明:

状态的存储和容错:由于状态是每个键和窗口的,因此您希望同时处理的键和窗口越多,您将产生的存储就越多。

但是,我找不到更多关于此的信息。

【问题讨论】:

    标签: apache-beam


    【解决方案1】:

    对于您问题的第一部分,beam 通过耗尽处理流服务中的异常,这里介绍了一些细节https://cloud.google.com/dataflow/docs/guides/stopping-a-pipeline

    虽然不确定这是否回答了您关于偏移量/检查点的问题。

    【讨论】:

    • 正如我目前所知,Apache Beam 中没有内置容错功能来存储在 Python 中的 DirectRunner、Flink 或 DataFlow Runner 上运行的全局窗口的状态。 Stopping A Pipeline 文章指出:不支持排空使用 Apache Beam SDK for Python 的流式传输管道。不支持排空批处理管道。我希望尽快解决这个问题。
    • Beam 已经添加了尚未包含在 Beam 模型中的其他常见功能。 Drain 和 Checkpoint 支持的文档还没有真正支持。beam.apache.org/documentation/runners/capability-matrix/…
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-08-28
    • 2018-11-10
    • 2017-01-04
    • 1970-01-01
    • 2021-10-20
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多