Flink的广播状态的作用有点类似hadoop hive的mapjoin,它将吞吐较小的一个流的数据Broadcast到下游的每个Task中,使得这些数据记录能够为所有的Task所共享,这样可以达到节省内存的效果.因此,广播变量比较适用的场景在动态数据的更新计算.
比如,对流上的数据做规则校验,符合某种规则的会被拦截下来,但是通常校验的规则会改变,如果因为改变校验规则又要停下应用去更新规则然后再去计算,这样的话就不太友好了.
这里来构造一个场景看看Broadcast State的效果:
对每条接收到的数据流进行关键字检测(假设初始化的关键字是Flink),如果关键字符合规则配置的则把这条消息流拦截下来,并且规则配置会定期的更新.
具体的实现如下:




应用输出结果:

从作业输出结果来看,当数据流发送消息包含Flink关键字时候,会被拦截下来,并且配置的规则是可以自动根据需要来更新,而无需停止应用.另外,这里设置了作业的并发数是4,消息在4个并发分区都能读到配置流的数据.
追踪Flink的源码来看,在实现Broadcast state时候,有以下几个方面需要注意:
首先需要创建MapStateDescriptor对象,而MapStateDescriptor的底层是Map结构;
通过DataStream.broadcast方法返回BroadcastStream;
主数据流跟BroadcastStream进行connect后来处理数据,这里会调用Process方法来处理.process接收两种类型的function,一种是KeyedBroadcastProcessFunction适用于KeyedStream,另外一种是
BroadcastProcessFunction.
BroadcastProcessFunction和KeyedBroadcastProcessFunction都包含了processElement和processBroadcastElement这两个抽象方法.




