【问题标题】:Flink Broadcast StateFlink 广播状态
【发布时间】:2021-09-07 01:43:02
【问题描述】:
  • 我有一个 RichParallelSourceFunction(parallelism=1),它每 5 秒查询一次 MySQL,它发出一个时间戳列表,指示何时开始/停止写入接收器。
  • 此时间戳广播到原始流(并行度=10)。
  • 我已将 RichParallelSourceFunction 的并行度设置为 1,以减少对 MySQL 服务器的同时请求数。
  • 我很困惑在这种情况下是否需要广播状态。为什么不将广播数据存储在运营商本地数据结构中?
  • .broadcast(stateDesc) 与 .broadcast() 之间有什么区别?
class MyBroadcastProcessFunction(name: String) extends BroadcastProcessFunction[Log, TimestampList, Log] with CheckpointedFunction {
    
      private var sortedTimestamps: IndexedSeq[Long] = _
      @transient var buffer: ListState[(Log, Long)] = _
    
      def processElement(value: Log, timestamp: Long, out: Collector[Log]): Unit = {
        if (timestamp > sortedTimestamps(0))
          out.collect(value)
        else
          buffer.put(value, 0L)
      }
    
    
      override def processBroadcastElement(timestamps: TimestampList,
                                           ctx: BroadcastProcessFunction[Log, TimestampList, Log]#Context,
                                           out: Collector[Log]): Unit = {
        // add timestamps to SortedTimestamps
        // do I need to use BroadcastState mapstate here? or just use operator-local data structures (ex. sortedTimestamps)?
    }
    
    override def snapshotState(context: FunctionSnapshotContext): Unit = {}
    
    override def initializeState(context: FunctionInitializationContext): Unit = {
        val stateDesc = new ListStateDescriptor[(Log, Long)]("logBuffer",
          classOf[(Log, Long)])
        buffer = context.getOperatorStateStore.getListState(stateDesc) 
    }
}

【问题讨论】:

    标签: scala apache-flink


    【解决方案1】:

    无论流的并行度如何,转换.broadcast() 都会将所有事件发送给所有下游操作符。 Doc says:

    设置 DataStream 的分区,以便将输出元素广播到下一个操作的每个并行实例。 返回:

    .broadcast(stateDesc) 是定义一个pattern state,您可以在其中找到一个基于另一个通常非常小的流的事件流的模式。 This also is a good reference.

    您创建BroadcastProcessFunction 的方式是错误的,因为您只处理一个流。处理广播状态的正确方法,在你的情况下是来自 MySql 的时间戳,在 processBroadcastElement() 方法。在这种方法中,您必须更新全局/广播状态。

    然后另一种方法processElement() 您会收到一个常规或快速流,您可以在其中找到基于您在第一种方法processBroadcastElement() 上更新的状态的模式。

    以下是您应该如何实施的更多信息。有一些注意事项,例如您将无法更新 ListState。最好使用链接中描述的MapState。

         def processBroadcastElement(timestamps: TimestampList,
                                               ctx: BroadcastProcessFunction[Log, TimestampList, Log]#Context, out: Collector[Log]): Unit = {
               // update buffer state
               // I don't think you can use .put() to update the ListState.
               // Actually I think it is not possible to update a ListState, than you have to use MapState.
               context.getOperatorStateStore.getListState(stateDesc)
                   .put(value.name, value);
          }
          def processElement(value: Log, timestamp: Long, out: Collector[Log]): Unit = {
            buffer = context.getOperatorStateStore.getListState(stateDesc)
            if (value: Log match within buffer ?)
              out.collect(value) // MATCH
          }
    

    【讨论】:

      【解决方案2】:

      .broadcast(stateDesc) 需要一个状态描述符,以便它知道如何序列化正在广播的数据。无论您是否希望将 MapState 中广播的数据存储在 KeyedBroadcastProcessFunction 中,您都可以使用它。

      如果您不使用 MapState 来存储这些数据,那么如果作业失败并重新启动,这些数据将会丢失。但也许这对您来说无关紧要,因为您可以在作业重新启动时从 MySQL 获取最新数据。

      【讨论】:

      • 你好大卫。我需要在作业重新启动时根据恢复的广播状态设置一些变量。我可以在 KeyedBroadcastProcessFunction 的 open() 中执行“getRuntimeContext.getBroadCastState(...)”吗?
      • @objectt 我相信是这样,但我还没有尝试过。这个想法本身并没有错,afaik。
      猜你喜欢
      • 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
      相关资源
      最近更新 更多