【发布时间】:2018-09-12 07:49:20
【问题描述】:
我得到一个表单的数据流:
+--+---------+---+----+
|id|timestamp|val|xxx |
+--+---------+---+----+
|1 |12:15:25 | 50| 1 |
|2 |12:15:25 | 30| 1 |
|3 |12:15:26 | 30| 2 |
|4 |12:15:27 | 50| 2 |
|5 |12:15:27 | 30| 3 |
|6 |12:15:27 | 60| 4 |
|7 |12:15:28 | 50| 5 |
|8 |12:15:30 | 60| 5 |
|9 |12:15:31 | 30| 6 |
|. |... |...|... |
我有兴趣将窗口操作应用于xxx 列,就像 Spark Streaming 中提供的时间戳窗口操作一样,具有一些窗口大小和滑动步长。
让groupBy下面有窗口函数,lines代表一个流数据帧,窗口大小:2,滑动步长:1。
val c_windowed_count = lines.groupBy(
window($"xxx", "2", "1"), $"val").count().orderBy("xxx")
所以,输出应该如下:
+------+---+-----+
|window|val|count|
+------+---+-----+
|[1, 3]|50 | 2 |
|[1, 3]|30 | 2 |
|[2, 4]|30 | 2 |
|[2, 4]|50 | 1 |
|[3, 5]|30 | 1 |
|[3, 5]|60 | 1 |
|[4, 6]|60 | 2 |
|[4, 6]|50 | 1 |
|... |.. | .. |
我尝试使用partitionBy,但 Spark Structured Streaming 不支持它。
我正在使用 Spark Structured Streaming 2.3.1。
谢谢!
【问题讨论】:
标签: scala apache-spark spark-streaming aggregate-functions spark-structured-streaming