【问题标题】:Flink: ProcessWindowFunctionFlink:进程窗口函数
【发布时间】:2017-10-04 03:30:02
【问题描述】:

我最近在研究 Flink 新版本中的ProcessWindowFunction。它说ProcessWindowFunction 支持全局状态和窗口状态。我使用 Scala API 来试一试。到目前为止,我可以让全局状态正常工作,但我没有任何运气让它成为窗口状态。我正在做的是处理系统日志并计算由主机名和严重级别键入的日志数量。我想计算两个相邻窗口之间的日志计数差异。这是我实现ProcessWindowFunction的代码。

class LogProcWindowFunction extends ProcessWindowFunction[LogEvent, LogEvent, Tuple, TimeWindow] {
  // Create a descriptor for ValueState
  private final val valueStateWindowDesc = new ValueStateDescriptor[Long](
    "windowCounters",
    createTypeInformation[Long])

  private final val reducingStateGlobalDesc = new ReducingStateDescriptor[Long](
    "globalCounters",
    new SumReduceFunction(),
    createTypeInformation[Long])

  override def process(key: Tuple, context: Context, elements: Iterable[LogEvent], out: Collector[LogEvent]): Unit = {
    // Initialize the per-key and per-window ValueState
    val valueWindowState = context.windowState.getState(valueStateWindowDesc)
    val reducingGlobalState = context.globalState.getReducingState(reducingStateGlobalDesc)
    val latestWindowCount = valueWindowState.value()
    println(s"lastWindowCount: $latestWindowCount ......")
    val latestGlobalCount = if (reducingGlobalState.get() == null) 0L else reducingGlobalState.get()
    // Compute the necessary statistics and determine if we should launch an alarm
    val eventCount = elements.size
    // Update the related state
    valueWindowState.update(eventCount.toLong)
    reducingGlobalState.add(eventCount.toLong)
    for (elem <- elements) {
      out.collect(elem)
    }
  }
}

我总是从窗口状态获得0 值,而不是之前更新的计数。我已经为这个问题苦苦挣扎了好几天。有人可以帮我弄清楚吗?谢谢。

【问题讨论】:

    标签: scala apache-flink flink-streaming


    【解决方案1】:

    每个窗口状态的范围是单个窗口实例。对于上面的 process 方法,每次调用它时都会在范围内创建一个新窗口,因此 latestWindowCount 始终为零。

    对于一个只会触发一次的普通普通窗口,每个窗口的状态是无用的。只有当一个窗口以某种方式有多次触发(例如,延迟触发)时,您才能充分利用每个窗口的状态。如果您想从一个窗口到另一个窗口记住某些内容,那么您可以使用全局窗口状态来做到这一点。

    有关使用每个窗口状态来记住要在后期触发中使用的数据的示例,请参阅 Flink 的 advanced window training 中的幻灯片 13-19。

    【讨论】:

    • 感谢您的回复。实际上,我发布的代码是从您提供的 URL 中模仿的。我的想法是我们可以启动每个窗口的状态,对这些状态进行一些计算,然后最后更新它们,以便我们可以在下一个窗口中检索它们。因此,我希望在对ProcessWindowFunction 的最新调用中保留窗口状态的值。我很困惑为什么即使我在调用ProcessWindowFunction 时更新了相应的窗口状态,你也会期望latestWindowCount 为零。我误解了每个窗口状态的用法吗??
    • 我已经扩展了我的答案,希望更清楚。如果您还有问题,请告诉我。基本上归结为“下一个窗口”不是同一个窗口,并且无法访问前一个窗口的每个窗口状态。
    • 感谢您更新答案。确实,我误解了每个窗口状态的定义。现在我可以完全弄清楚了。
    猜你喜欢
    • 2020-07-16
    • 1970-01-01
    • 2020-08-30
    • 1970-01-01
    • 2017-10-17
    • 2017-08-03
    • 2017-05-25
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多