【发布时间】: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