Flink 1.13.1中两种Broadcast State使用方式的正确性、性能及内存占用对比分析
关于Flink Broadcast State两种使用方式的正确性与性能分析
先给大家梳理下背景:我们用的是Flink 1.13.1版本,项目里用到了Broadcast State,每5分钟会拉取一批配置;因为一个处理函数只能关联一个广播源,所以定义了一个Case Class来封装三类配置,代码如下:
case class OrderConfBroadcastBean( orderRuleConfig: List[OrderInfoBean], userSegmentInfo: Map[String, (String, String)], lacciRegRel: Map[String, Set[String]] )
广播状态的初始化代码是这样的:
val orderConfBroadcast = env.addSource(new OrderConfSource(dbConfig, serverConfig.smsRuleRedis)) .name("order_conf_load") .uid("order_conf_load") .setParallelism(1) .broadcast(new MapStateDescriptor[String, OrderConfBroadcastBean]( "order_conf_broadcast", createTypeInformation[String], createTypeInformation[OrderConfBroadcastBean] ))
现在在处理函数里有两种使用Broadcast State的写法,想请教下:哪种方式是正确的?或者哪种具备更优的性能与更低的内存占用?原因是什么?
第一种实现方式
class OrderFilterProcess( var userSegmentInfo: Map[String, (String, String)], var orderInfo: List[OrderInfoBean], redisConf: String, var lacciRegRel: Map[String, Set[String]] ) extends KeyedBroadcastProcessFunction[String, RegLacciBean, OrderConfBroadcastBean, OrderResultBean] { override def processElement( regLacci: RegLacciBean, ctx: KeyedBroadcastProcessFunction[String, RegLacciBean, OrderConfBroadcastBean, OrderResultBean]#ReadOnlyContext, out: Collector[OrderResultBean] ): Unit = { userSegmentInfo.get("xxx") orderInfo.map(xxx) } override def processBroadcastElement( value: OrderConfBroadcastBean, ctx: KeyedBroadcastProcessFunction[String, RegLacciBean, OrderConfBroadcastBean, OrderResultBean]#Context, out: Collector[OrderResultBean] ): Unit = { if (value.orderRuleConfig.nonEmpty) { orderInfo = value.orderRuleConfig } if (value.userSegmentInfo.nonEmpty) { userSegmentInfo = value.userSegmentInfo } if (value.lacciRegRel.nonEmpty) { lacciRegRel = value.lacciRegRel } } }
第二种实现方式
class OrderFilterProcess( var userSegmentInfo: Map[String, (String, String)], var orderInfo: List[OrderInfoBean], redisConf: String, var lacciRegRel: Map[String, Set[String]] ) extends KeyedBroadcastProcessFunction[String, RegLacciBean, OrderConfBroadcastBean, OrderResultBean] { val stateDescriptor = new MapStateDescriptor[String, OrderConfBroadcastBean]( "order_conf_broadcast", createTypeInformation[String], createTypeInformation[OrderConfBroadcastBean] ) override def processElement( regLacci: RegLacciBean, ctx: KeyedBroadcastProcessFunction[String, RegLacciBean, OrderConfBroadcastBean, OrderResultBean]#ReadOnlyContext, out: Collector[OrderResultBean] ): Unit = { val state = ctx.getBroadcastState(stateDescriptor) Option(state.get("order_state")).map(_.userSegmentInfo.get("xxx")).orElse(userSegmentInfo.get("xxx")) } override def processBroadcastElement( value: OrderConfBroadcastBean, ctx: KeyedBroadcastProcessFunction[String, RegLacciBean, OrderConfBroadcastBean, OrderResultBean]#Context, out: Collector[OrderResultBean] ): Unit = { ctx.getBroadcastState(stateDescriptor).put("order_state", value); } }
正确性与性能内存分析
咱们先直接给出结论:第二种方式是正确的生产级实现,第一种方式存在致命的容错问题,绝对不能用于生产环境,具体原因如下:
1. 正确性层面
- 第一种方式的核心问题:它把广播配置存在了处理函数的成员变量里,完全绕开了Flink的状态管理机制。Flink的状态需要参与checkpoint/savepoint的持久化,一旦算子实例重启、故障恢复或者集群扩缩容,这些成员变量里的配置会丢失,导致计算逻辑完全错误。而且广播源的更新只会发送到每个并行实例的
processBroadcastElement方法,成员变量的更新是每个实例各自维护的,无法保证集群内所有并行实例的配置一致性(比如实例重启后,成员变量会回到初始化值)。 - 第二种方式是Flink官方推荐的Broadcast State使用方式:通过
ctx.getBroadcastState获取的状态会被Flink统一管理,自动参与checkpoint,故障恢复时能从快照中恢复最新的配置;同时Flink会保证广播源的更新会被分发到所有并行实例的Broadcast State中,确保所有实例使用的配置是一致的。
2. 性能与内存层面
- 内存占用:两种方式其实差不多,因为广播状态的特性就是每个并行实例都会持有一份完整的广播数据副本(不管是存在成员变量还是Broadcast State里)。但第二种方式的内存使用是被Flink监控和管理的,能通过Flink的状态指标查看内存占用情况,而第一种方式的内存占用无法被Flink追踪。
- 性能:第一种方式直接访问成员变量看起来会快一点,但这种“性能优势”完全是建立在牺牲容错的基础上的,没有任何生产价值。第二种方式的Broadcast State是Flink优化过的状态存储,访问性能足够满足绝大多数场景,而且能保证数据的一致性和可靠性。
另外补充一点:第二种方式里,Flink会保证同一并行实例中processElement和processBroadcastElement是串行执行的,所以不会出现并发访问Broadcast State的问题,不需要额外加锁,线程安全有保障。
内容的提问来源于stack exchange,提问作者Shigure
相关产品推荐
相关产品推荐

