【问题标题】:Flink stateful function address resolution for messaging用于消息传递的 Flink 有状态函数地址解析
【发布时间】:2020-12-12 07:49:51
【问题描述】:

在Flink数据流中假设上游算子托管在机器/任务管理器m上,上游算子如何知道下游算子所在的机器(任务管理器)m’。是不是在JobManager对作业子/任务(算子)的初始调度过程中建立了这样的下游/上游算子之间的数据流路径,并且这些数据流路径在应用生命周期内是固定的?

更一般地,考虑支持动态消息传递且数据流未固定或未预定义的 Flink 有状态函数,并给定一个带有键 k 的函数,该函数需要将消息/事件发送到带有键 k’ 的另一个函数函数k 如何找到函数k’ 的地址来发送消息? Flink 运行时是否在某些分布式数据结构(例如 Microsoft Orleans 中的 DHT)中保留键机映射,并且每次调用函数都涉及对此类数据结构的访问?

请注意,我来自 Spark 背景,在给定 RDD/batch 模型的情况下,作业图任务是连续执行的(在 shuffle 边界处中断),并且每个 shuffle 子任务都被指示持有应该被拉出的键子集的机器/由该子任务处理...。

谢谢。

【问题讨论】:

    标签: apache-flink partitioning actor flink-streaming flink-statefun


    【解决方案1】:

    即使是有状态的函数,底层 Flink 作业的拓扑结构在作业启动时也是固定的。每个有状态函数作业都或多或少地使用这样的作业图(入口各不相同,但其余的总是这样):

    这里你看到所有加载的入口都变成了 Flink 源操作符,发出输入消息, 并且路由器成为链接到这些来源的平面地图操作员。

    充当路由器的平面图将输入消息转换为内部事件信封,从而 本质上只是将消息有效负载与其目标逻辑地址包装在一起。信封是 流经流图的所有消息的在线数据类型。 Stateful Functions 运行时以函数调度器运算符为中心, 它跨所有模块运行所有已加载函数的实例。

    在router flatmap operator和function dispatcher operator之间是一个keyBy操作 它使用目标目标id 作为键对输入流进行重新分区。这 network shuffle 保证所有发往给定id 的消息都发送到同一个 函数调度运算符的实例。

    收到后,函数调度器从信封中提取目标函数地址,加载 该函数实例,然后使用包装的输入调用该函数(它也在 信封)。

    函数调度器的不同实例如何相互发送消息?

    这是通过将每个函数调度器与一个反馈运算符放在一起来完成的。 所有传出消息都使用目标函数id 作为键进行另一个网络洗牌。

    此反馈运算符在作业图中创建一个循环或迭代。有状态函数在其消息传递模式中可以有循环或循环,并且不限于使用 DAG 处理数据。

    反馈渠道已设置检查点;在失败的情况下消息永远不会丢失。

    关于这方面的更多信息,我推荐 Tzu-Li (Gordon) Tai 的 Flink Forward 演讲:Stateful Functions: Polyglot Event-Driven Functions for Stateful Distributed Applications。上图来自他的演讲。

    【讨论】:

      猜你喜欢
      • 2020-08-14
      • 1970-01-01
      • 1970-01-01
      • 2020-07-27
      • 2017-12-22
      • 1970-01-01
      • 2021-11-12
      • 2015-12-20
      • 2021-05-09
      相关资源
      最近更新 更多