【问题标题】:Apache Flink, second stage summing on a windowed streamApache Flink,窗口流的第二阶段求和
【发布时间】:2018-12-13 05:38:46
【问题描述】:

我有一个 Kafka 事件的数据流(对应于来自设备的读数)馈送到以下代码中,该代码在滑动窗口中生成每个设备的平均读数。这工作正常。接下来我想计算同一窗口中所有每个设备平均值的总和,这是我无法正确语法表达的部分。

这部分有效:

val stream = env
    // words is our Kafka topic
    .addSource(kafkaConsumer)
    // configure timestamp and watermark assigner
    .assignTimestampsAndWatermarks(new DeviceTSAssigner)
      .keyBy(_.deviceIdFull)
      .timeWindow(Time.minutes(5), Time.minutes(1))
    /* count events in window */
      .apply{ (key: String, window: TimeWindow, events: Iterable[DeviceData], out: Collector[(String, Long, Double)]) =>
        out.collect( (key, window.getEnd, events.map(_.currentReading).sum/events.size))
    }

  stream.print()

输出类似于

(device1,1530681420000,0.0)
(device2,1530681420000,0.0)
(device3,1530681480000,0.0)
(device4,1530681480000,0.0)
(device5,1530681480000,52066.0)
(device6,1530681480000,69039.0)
(device7,1530681480000,79939.0)
... 
...

以下代码是我遇到问题的部分,我不确定如何编码,但我认为它应该是这样的:

  val avgStream = stream
    .keyBy(2) // 2 represents the window.End from stream, see code above
    .timeWindow(Time.minutes(1)) // tumbling window
    .apply { (
               key: Long,
               window: TimeWindow,
               events: Iterable[(String, Long, Double)],
               out: Collector[(Long, Double)]) =>
      out.collect( (key, events.map( _._3 ).sum ))
    }

编译此代码时出现以下错误..

Error:(70, 52) type mismatch;
 found   : (Long, org.apache.flink.streaming.api.windowing.windows.TimeWindow, Iterable[(String, Long, Double)], org.apache.flink.util.Collector[(Long, Double)]) => Unit
 required: (org.apache.flink.api.java.tuple.Tuple, org.apache.flink.streaming.api.windowing.windows.TimeWindow, Iterable[(String, Long, Double)], org.apache.flink.util.Collector[?]) => Unit
                   out: Collector[(Long, Double)]) =>

我也尝试了其他变体,例如使用 AggregtionFunctions,但无法通过编译。似乎我需要将输入流元素转换为元组的错误,我已经看到了一些代码,但不完全确定如何做到这一点。我是 Scala 的新手,所以我认为这是这里的主要问题,因为我想做的并不复杂。

2018 年 7 月 4 日更新

我认为我的问题有一个解决方案,似乎工作正常,但我仍然希望保持开放,希望其他人可以评论它(问题以及我的解决方案)。

基本上,我删除了第一个字段(通过映射),它是设备的名称,因为我们不需要它,然后是时间戳上的 keyBy(来自前一阶段),将事件窗口化到一个翻滚窗口,然后只需对第二个(基于索引 1、0)字段求和,这是前一阶段的平均读数。

val avgStream = stream
          .map(r => (r._2, r._3))
         .keyBy(0)
         .timeWindowAll(Time.minutes(1))
         .sum(1)
         .print()

【问题讨论】:

  • 我想你的意思是 .keyBy(1) 如果你想使用第二个字段作为你的记录作为键。索引从零开始。
  • 映射流只有两个字段,即(时间戳,每个设备的平均读数),因此 keyBy 第一个字段的索引将为 0。

标签: scala apache-flink flink-streaming


【解决方案1】:

我能够回答我自己的问题,因此上述方法(请参阅 2018 年 7 月 4 日更新)有效,但更好的方法是执行此操作(尤其是如果您不想只针对stream but multiple) 是使用 AggregateFunction。我之前也尝试过,但由于缺少“地图”步骤而遇到了问题。

一旦我在第二阶段映射流以提取相关的感兴趣领域,我就可以使用 AggregateFunction。

Flink 文档 heregithub link 都为此提供了示例。我从 Flink 文档示例开始,因为它非常易于理解,然后将我的代码转换为看起来更像 github 示例。

【讨论】:

    猜你喜欢
    • 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
    相关资源
    最近更新 更多