【问题标题】:Pyspark structured streaming window (moving average) over last N data points最后 N 个数据点的 Pyspark 结构化流窗口(移动平均)
【发布时间】:2020-05-09 16:17:25
【问题描述】:

我使用 Pyspark Structured Streaming 2.4.4 从 kafka 主题中读取了几个数据帧。我想向该数据帧添加一些新列,这些新列主要基于过去 N 个数据点的窗口计算(例如:过去 20 个数据点的移动平均值),并且随着新数据点的交付,相应的值MA_20 应立即计算。

数据可能如下所示: 时间戳 |波动率指数

2020-01-22 10:20:32 | 13.05
2020-01-22 10:25:31 | 14.35
2020-01-23 09:00:20 | 14.12

值得一提的是,数据将在周一至周五每天 8 小时内收到。 因此周一早上计算的移动平均线应该包括周五的数据!

我尝试了不同的方法,但仍然无法实现我想要的。

windows = df_vix \
    .withWatermark("Timestamp", "100 minutes") \
    .groupBy(F.window("Timestamp", "100 minute", "5 minute")) \

aggregatedDF = windows.agg(F.avg("VIX"))

前面的代码计算了 MA,但它会将周五的数据视为较晚,因此将它们排除在外。比最后 100 分钟更好应该是最后 20 分(间隔 5 分钟)。

我认为我可以使用 rowsBetween 或 rangeBetween,但在流数据帧中,窗口不能应用于非时间戳列 (F.col('Timestamp').cast('long'))

    w = Window.orderBy(F.col('Timestamp').cast('long')).rowsBetween(-600, 0)

    df = df_vix.withColumn('MA_20', F.avg('VIX').over(w)

)

但另一方面,不可能在 rowsBetween() 中指定间隔,使用 rowsBetween(- minutes(20), 0) 抛出:分钟未定义(sql.functions 中没有这样的函数)

我找到了另一种方式,但它也不适用于流式数据帧。不知道为什么会出现“流数据帧不支持非基于时间的窗口”错误(df_vix.Timestamp 属于时间戳类型)

df.createOrReplaceTempView("df_vix")

df_vix.createOrReplaceTempView("df_vix")
aggregatedDF = spark.sql(
    """SELECT *, mean(VIX) OVER (
        ORDER BY CAST(df_vix.Timestamp AS timestamp)
        RANGE BETWEEN INTERVAL 100 MINUTES PRECEDING AND CURRENT ROW
     ) AS mean FROM df_vix""")

我不知道我还能用什么来计算简单的移动平均线。看起来在 Pyspark 中实现这一点是不可能的......也许更好的解决方案是每次新数据将整个 Spark 数据帧传递给 Pandas 并计算 Pandas 中的所有内容(或将新行附加到 Pandas 并计算 MA)时进行转换? ??

我认为随着新数据的到来创建新功能是结构化流的主要目的,但事实证明 Pyspark 不适合这个,我正在考虑放弃 Pyspark 转而使用 Pandas ...

编辑

尽管 df_vix.Timestamp 类型为“timestamp”,但以下内容也无法正常工作,但它会引发“流数据帧不支持非基于时间的窗口”错误。

w = Window.orderBy(df_vix.Timestamp).rowsBetween(-20, -1)
aggregatedDF = df_vix.withColumn("MA", F.avg("VIX").over(w))

【问题讨论】:

    标签: python apache-spark pyspark spark-streaming


    【解决方案1】:

    你看过window在活动时间里的操作吗? window(timestamp, "10 minutes", "5 minutes") 将每 5 分钟为您提供 10 分钟的数据框,然后您可以对其进行聚合,包括移动平均线。

    【讨论】:

    • 这在多台设备发送数据的情况下不起作用,并且某些设备比当前时间稍晚,窗口功能将简单地忽略落后5分钟的设备的数据。跨度>
    猜你喜欢
    • 2021-07-07
    • 2022-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-09-24
    • 2020-03-29
    • 2016-05-26
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多