首页 > 解决方案 > Flink WindowFunction 折叠

问题描述

我创建了一个滑动窗口并希望递归地打包所有进入该窗口期间的元素,这是代码的一部分

.map(x => ((x.pickup.get.latitude, x.pickup.get.longitude), (x.dropoff.get.latitude, x.dropoff.get.longitude)))
        .windowAll(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1)))
        .fold(List[((Double, Double), (Double, Double))]) {(acc, v) => acc :+ ((v._1._1, v._1._2), (v._2._1, v._2._2))}

我希望创建一个List其中的元素tuple,但这不起作用。

我试过这个并且它有效:

val l2 : List[((Int, Int), (Int, Int))] = List(((1, 1), (2, 2)))
val newl2 = l2 :+ ((3, 3), (4, 4))

我怎样才能做到这一点?非常感谢

标签: scalaapache-flinkflink-streaming

解决方案


函数的第一个参数fold需要是初始值而不是类型。将最后一行更改为:

.fold(List.empty[((Long, Long), (Long, Long))]) {(acc, v) => acc :+ ((v._1._1, v._1._2), (v._2._1, v._2._2))}

应该做的伎俩。


推荐阅读