【问题标题】:Flink stateful functions : compensating callback on a timeoutFlink 有状态函数:超时补偿回调
【发布时间】:2020-09-06 14:39:02
【问题描述】:

我正在 Flink 有状态函数中实现一个用例。我的规范强调,从一个有状态函数 f 开始,一个业务工作流(换句话说,一组有状态函数 f1、f2、... fn 被顺序或并行调用,或者两个都 )。 有状态函数 f 等待返回结果以更新本地状态,它还启动超时回调,即给自己的消息。在超时时,f 检查本地状态是否已更新(它已收到结果),如果是这样的话,生命是好的。

但是,如果在超时时 f 发现它还没有收到结果,它必须启动补偿工作流来撤消 有状态函数 f1、f2、... fn 可能收到的任何更改。

Flink 有状态功能框架是否支持设计模式/用例等,还是应该在应用程序级别实现?实现这种解决方案的最简单设计是什么?例如,如何知道工作流有状态函数 f1, f2, ... fn 的哪些函数受到超时调用(控制流已超时)的影响? Flink sateful 函数和集成消息和状态的概念如何促成这种模式?

谢谢。

【问题讨论】:

    标签: apache-flink flink-statefun


    【解决方案1】:

    我在 Apache Flink 邮件列表上发布了这个问题,并得到了 Igal Shilman 的以下回复,感谢 Igal。

    我想首先提到的是,如果您的原始 这种情况的动机是对短暂故障的关注,例如:

    • 函数 Y 是否曾收到函数 X 发送的消息?
    • 发送消息失败了吗?
    • 目标函数是否可以接受发送给它的消息?
    • 消息的顺序是不是搞错了?
    • 等'

    然后,StateFun 消除了所有这些问题和一整类 临时错误,否则您将不得不自己处理 您的业​​务逻辑(如重试、退避、服务发现等)。

    现在,如果您的激励方案不是关于暂时性错误,而是更多 关于事务性工作流程,然后正如 Dawid 提到的,您将不得不 实施 这在您的应用程序逻辑中。我认为你描述的方式 流应该直接映射到一个协调函数(每个流实例) 跟踪其内部状态的结果/超时。

    这是一个草图:

    1. 流协调器函数 - 将使用输入调用 需要启动流程。它将开始调用相关的 函数(由流的 DAG 定义)并保持内部状态 表示 调用了哪些函数(地址)及其完成状态。 当流程成功完成时,协调器可以安全地丢弃它的 状态。 在协调器决定中止流程的任何情况下(内部 超时/外部消息/等')它必须检查其内部 状态并启动补偿工作流程(向 已经成功/正在进行的函数)

    2. 流程中的每个函数都必须接受来自协调器的消息, 依次进行,并回复成功或失败。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-07-27
      • 2020-08-17
      • 1970-01-01
      • 2022-11-30
      • 1970-01-01
      • 2020-08-29
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多