【问题标题】:Stateful Functions Fault Tolerant Message Distribution in Apache FlinkApache Flink 中的有状态函数容错消息分发
【发布时间】:2020-08-14 16:58:53
【问题描述】:

我正在尝试使用 apache flink 有状态函数来实现消息传递场景。

根据设计,我需要从传入消息中计算一些统计数据并将它们存储在状态中。之后,场景函数将访问这些状态和消息并在它们上运行业务规则。但是我们每条消息可能有几十个场景,每个场景都应该只运行一次。

代码大致如下

@Override
    public void configure(MatchBinder binder) {
        binder
            .predicate(Transaction.class,this::updateTransactionStatAndSendToScenatioManager)
}

    private void updateTransactionStatAndSendToScenatioManager(Context context, Transaction transaction){
        // state update
        context.send(FnScenarioManager.TYPE,  String.valueOf(transaction.id()) , transaction);
    }

FnScenarioManager:

@Override
    public void configure(MatchBinder binder) {
    binder
        .predicate(Transaction.class,this::runTransactionScenarios);
}


private void runTransactionScenarios(Context context, Transaction transaction){
   context.send(Scenario1.TYPE,String.valueOf(transaction.id()),transaction);
   context.send(Scenario2.TYPE,String.valueOf(transaction.id()),transaction);
   context.send(Scenario3.TYPE,String.valueOf(transaction.id()),transaction);
   ...
   context.send(ScenarioN.TYPE,String.valueOf(transaction.id()),transaction);
}

我的问题是如果集群在 runTransactionScenarios 中间崩溃会发生什么?

  • 每个场景会只运行一次吗?如果不是,我该如何确保?

【问题讨论】:

    标签: state apache-flink stateful flink-statefun


    【解决方案1】:

    有状态函数(以及一般的 Apache Flink)支持完全一次的状态语义。这意味着在发生故障的情况下,运行时将始终以模拟完全无故障执行的方式回滚状态和消息。

    这意味着消息可能会被重播,但内部状态将回滚到收到消息之前的时间点。只要您的业务规则仅修改 statefun 状态并通过出口与外界交互,您就可以将系统视为具有仅一次的属性。

    【讨论】:

    • 如果场景 1 更新内部状态然后集群在其他场景更新其内部状态之前失败会发生什么?如果我们使用 state 作为计数器,它会更新两次,对吗?
    • 否,在这种情况下,场景 1 会将其状态回滚到消息之前的值,然后重新处理它,从而产生相同的最终状态。虽然它在技术上会两次处理消息,但不会导致错误的结果。想象一下您只是在进行计数,当前计数为 5,并且一条消息将计数增加到 6。如果发生故障,计数将重置为 5,然后重新处理该消息,从而产生相同的正确结果。
    • 假设我们有 5 个按顺序工作的计数场景作为我最初的问题。场景 1,2,3 增加了他们的状态 4,5 没有。在失败的情况下,flink 如何为每个场景计算正确的场景状态。是否有关于在失败时状态如何表现的文档?
    猜你喜欢
    • 2020-07-27
    • 1970-01-01
    • 2020-08-29
    • 1970-01-01
    • 2020-12-12
    • 1970-01-01
    • 1970-01-01
    • 2020-05-17
    • 1970-01-01
    相关资源
    最近更新 更多