Flink WindowFunction折叠

问题描述 投票:0回答:1

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

.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))

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

scala apache-flink flink-streaming
1个回答
0
投票

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

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

应该做的伎俩。

© www.soinside.com 2019 - 2024. All rights reserved.