使用 SCALA 滑动 vs mllib 滑动 - 两种实现,有点繁琐,但这里是:
import org.apache.spark.mllib.rdd.RDDFunctions._
val rdd1 = sc.parallelize(Seq(
( "key1", "value1"),
( "key2", "value2"),
( "key3", "value3"),
( "key4", "value4"),
( "key5", "value5")
))
val rdd2 = rdd1.sliding(2)
val rdd3 = rdd2.map(x => (x(0), x(1)))
val rdd4 = rdd3.map(x => ((x._1._1, x._2._1),x._1._2, x._2._2))
rdd4.collect
还有,下面这个当然更好……:
val rdd5 = rdd2.map{case Array(x,y) => ((x._1, y._1), x._2, y._2)}
rdd5.collect
两种情况都返回:
res70: Array[((String, String), String, String)] = Array(((key1,key2),value1,value2), ((key2,key3),value2,value3), ((key3,key4),value3,value4), ((key4,key5),value4,value5))
我相信它可以满足您的需求,但在 pyspark 中不满足。
在 Stack Overflow 上,您可以找到 pyspark 没有 RDD 等价物的声明,除非您“自己动手”。你可以看看这个How to transform data with sliding window over time series data in Pyspark。但是,我会建议使用 pyspark.sql.functions.lead() 和 pyspark.sql.functions.lag() 的数据框。稍微容易一些。