我在想怎么写我自己的
WindowFunction
但是有问题,我不知道为什么。我遇到的问题是apply函数,因为它无法识别
MyWindowFunction
作为有效输入,所以我无法编译。我正在传输的数据包含
(timestamp,x,y)
其中,对于测试,x和y分别为0和1。
extractTupleWithoutTs
只返回一个元组
(x,y)
. 我已经成功地用简单的求和和和归约函数运行了代码。感谢您的帮助:)使用Flink 1.3
进口:
import org.apache.flink.streaming.api.TimeCharacteristic
import org.apache.flink.streaming.api.functions.AssignerWithPeriodicWatermarks
import org.apache.flink.streaming.api.scala.function.WindowFunction
import org.apache.flink.streaming.api.watermark.Watermark
import org.apache.flink.streaming.api.windowing.windows.TimeWindow
import org.apache.flink.util.Collector
其余代码:
val env = StreamExecutionEnvironment.getExecutionEnvironment
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)
val text = env.socketTextStream("localhost", 9999).assignTimestampsAndWatermarks(new TsExtractor)
val tuple = text.map( str => extractTupleWithoutTs(str))
val counts = tuple.keyBy(0).timeWindow(Time.seconds(5)).apply(new MyWindowFunction())
counts.print()
env.execute("Window Stream")
MyWindow函数,基本上是从示例中复制粘贴并更改类型。
class MyWindowFunction extends WindowFunction[(Int, Int), Int, Int, TimeWindow] {
def apply(key: Int, window: TimeWindow, input: Iterable[(Int, Int)], out: Collector[Int]): () = {
var count = 0
for (in <- input) {
count = count + 1
}
out.collect(count)
}
}