【问题标题】:Which way of using flink's broadcast state is better使用flink的广播状态哪种方式更好
【发布时间】:2021-10-19 07:23:06
【问题描述】:

使用 flink 1.13.1 版本

代码已经精简。

我在我的项目中使用广播状态,它会每5分钟发送一些配置,因为一个进程函数只连接一个广播源,所以我定义了一个案例类来传输三种配置

案例类:

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]))

想知道在进程函数中使用广播状态的两种方式,哪一种是对的,或者哪一种性能更好,内存占用更低,为什么

第一次使用:

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(ruleStateDescriptor)
    Option(state.get("order_state")).map(_.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);
  }
}

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    这两种实现之间的最大区别在于,第一种是将广播流接收到的数据存储到作业失败时将丢失的变量中,而第二种实现是使用广播状态,该状态将被检查点并康复。

    第二版有一些开销。您必须对其进行测量才能确定有多少——但在这两种情况下,数据都将存储在内存中,因此差异应该不会很大。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多