【问题标题】:Spark Strucutured Streaming Window on non-timestamp column非时间戳列上的 Spark Structured Streaming 窗口
【发布时间】: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


    【解决方案1】:

    目前无法使用 Spark 结构化流以这种方式在非时间戳列上使用窗口。但是,您可以xxx 列转换为时间戳列,执行groupBycount,然后再转换回来。

    from_unixtime 可用于将自 1970-01-01 以来的秒数转换为时间戳。使用 xxx 列作为秒,并且可以创建一个 假时间戳 以在窗口中使用:

    lines.groupBy(window(from_unixtime($"xxx"), "2 seconds", "1 seconds"), $"val").count()
      .withColumn("window", struct(unix_timestamp($"window.start"), unix_timestamp($"window.end")).as("window"))
      .filter($"window.col1" =!= 0)
      .orderBy($"window.col1")
    

    上面,分组是在转换后的时间戳上完成的,下一行会将其转换回原来的数字。过滤器已完成,因为前两行将是一个窗口[0,2](即仅在具有xxx 等于1 的行上)但可以跳过。

    上述输入的结果输出:

    +------+---+-----+
    |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|
    | [5,7]| 30|    1|
    | [5,7]| 60|    1|
    | [5,7]| 50|    1|
    | [6,8]| 30|    1|
    +------+---+-----+
    

    【讨论】:

    • 非常感谢 Shaido!很棒的逻辑和完美的答案:)
    • @shaikh: 乐于助人:)
    【解决方案2】:

    spark 2.2 中的新功能是arbitrary-stateful-operations

    一个用例是管理用户会话,一个“用户窗口”

    scroll half way down this page to see an example

    如果 Shaido 的聪明解决方案对您有用,那么我建议您继续这样做。对于更复杂的需求,任意状态操作看起来像是要走的路。

    【讨论】:

      猜你喜欢
      • 2018-08-08
      • 1970-01-01
      • 2019-10-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-07-23
      • 1970-01-01
      • 2020-03-19
      相关资源
      最近更新 更多