【发布时间】: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