暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Flink流上的广播实现

Nathan的笔记 2020-04-11
1462

Flink的广播状态的作用有点类似hadoop hive的mapjoin,它将吞吐较小的一个流的数据Broadcast到下游的每个Task中,使得这些数据记录能够为所有的Task所共享,这样可以达到节省内存的效果.因此,广播变量比较适用的场景在动态数据的更新计算.

比如,对流上的数据做规则校验,符合某种规则的会被拦截下来,但是通常校验的规则会改变,如果因为改变校验规则又要停下应用去更新规则然后再去计算,这样的话就不太友好了.

这里来构造一个场景看看Broadcast State的效果:

  • 对每条接收到的数据流进行关键字检测(假设初始化的关键字是Flink),如果关键字符合规则配置的则把这条消息流拦截下来,并且规则配置会定期的更新.

具体的实现如下:



应用输出结果:


从作业输出结果来看,当数据流发送消息包含Flink关键字时候,会被拦截下来,并且配置的规则是可以自动根据需要来更新,而无需停止应用.另外,这里设置了作业的并发数是4,消息在4个并发分区都能读到配置流的数据.

追踪Flink的源码来看,在实现Broadcast state时候,有以下几个方面需要注意:

  1. 首先需要创建MapStateDescriptor对象,而MapStateDescriptor的底层是Map结构;

  2. 通过DataStream.broadcast方法返回BroadcastStream;

  3. 主数据流跟BroadcastStream进行connect后来处理数据,这里会调用Process方法来处理.process接收两种类型的function,一种是KeyedBroadcastProcessFunction适用于KeyedStream,另外一种是

    BroadcastProcessFunction.

  4. BroadcastProcessFunction和KeyedBroadcastProcessFunction都包含了processElement和processBroadcastElement这两个抽象方法.

文章转载自Nathan的笔记,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论