processElement()是否会阻塞broadcastProcessElement()?已知后者执行时前者不运行
关于KeyedBroadcastProcessFunction方法执行的调度问题
在Flink的KeyedBroadcastProcessFunction里,broadcastProcessElement()和processElement()的执行遵循单线程串行调度的规则,你的问题答案很明确:
- 已知
broadcastProcessElement()执行时processElement()不会运行,反过来,processElement()执行时同样会阻塞broadcastProcessElement(),二者无法并行执行。
核心原因是每个Flink算子任务都运行在单线程中,普通流和广播流的所有输入元素都会被放进同一个任务队列,按先到先得的顺序处理。不管是触发processElement()的普通元素,还是触发broadcastProcessElement()的广播元素,同一时间只能有一个方法在执行,另一个必须等当前方法执行完才能被调度。
举个实际场景的例子:
如果先收到普通流的元素,processElement()开始执行,这时候哪怕广播流有元素进来,也得等processElement()跑完,才会调用broadcastProcessElement()处理广播元素;反过来的情况也是一样的。
补充一点:这种串行限制只针对同一个算子实例,不同并行度的算子任务之间是各自独立运行的,互相不会干扰。
内容的提问来源于stack exchange,提问作者K.M
相关产品推荐
相关产品推荐

