【发布时间】: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