【问题标题】:Pyspark: re-sampling frequencies down to millisecondsPyspark:重新采样频率低至毫秒
【发布时间】:2019-10-19 11:18:48
【问题描述】:

我需要在时间戳列上加入两个火花数据帧。问题在于它们具有不同的频率:第一个数据帧 (df1) 每 10 分钟进行一次观察,而第二个数据帧 (df2) 是 25 hz(每秒 25 次观察,比 df1 频率高 15000 倍)。每个数据框有 100 多列和数百万行。为了进行平滑连接,我尝试将 df1 重新采样到 25 hz,预先填充由重新采样引起的 Null 值,然后在它们处于相同频率时加入数据帧。数据框太大,这就是我尝试使用 spark 而不是 pandas 的原因。

所以,问题来了:假设我有以下 spark 数据框:

我想将其重新采样到 25 赫兹(每秒 25 次观察),使其看起来像这样:

如何在 pyspark 中高效地做到这一点?

注意:

我尝试使用较早问题 (PySpark: how to resample frequencies) 中的代码重新采样我的 df1,如下所示:

from pyspark.sql.functions import col, max as max_, min as min_

freq = x   # x is the frequency in seconds

epoch = (col("timestamp").cast("bigint") / freq).cast("bigint") * freq 

with_epoch  = df1.withColumn("dummy", epoch)

min_epoch, max_epoch = with_epoch.select(min_("dummy"), max_("dummy")).first()

new_df = spark.range(min_epoch, max_epoch + 1, freq).toDF("dummy")

new_df.join(with_epoch, "dummy", "left").orderBy("dummy")
.withColumn("timestamp_resampled", col("dummy").cast("timestamp"))

看来,上述代码仅在预期频率大于或等于一秒时才有效。例如,当freq = 1时,它会产生下表:

但是,当我通过 25 hz 作为频率(即 freq = 1/25)时,代码会失败,因为 spark.range 函数中的“步长”不能小于 1。

是否有解决此问题的解决方法?或者任何其他方式将频率重新采样到毫秒?

【问题讨论】:

  • 我有一个类似的情况,而不是使用步长 = 1/25 的范围,您是否尝试过先将您的纪元转换为毫秒(乘以 1000)然后您的新步长为 40,并且然后,您可以使用:from_unixtime 将您的纪元转换回时间戳,格式为:“yyyy-MM-dd'T'HH:mm:ss.SSS”?在我的 SparkSQL 版本中,时间戳以秒为单位存储,但我有另一个后端捕获时间序列以纳秒为单位。要加入和重新采样,我必须先将后者重新缩放 10^9 倍,然后它才起作用。

标签: python pyspark resampling


【解决方案1】:

如果您的目标是连接 2 个数据框,我建议直接使用内部连接:

df = df1.join(df2, df1.Timestamp == df2.Timestamp)

但是,如果您想尝试对数据帧进行下采样,您可以将时间戳转换为毫秒并保留 mod(timestamp, 25) == 0 的那些行。只有当您确定数据被完美采样时,您才能使用它。

from pyspark.sql.functions import col
df1 = df1.filter( ((col("Timestamp") % 25) == 0 )

其他选项是对每一行进行编号并每 25 保留 1 个。使用此解决方案,您将减少行而不考虑时间戳。这个方案的另一个问题是需要对数据进行排序(效率不高)。

PD:过早的优化是万恶之源

编辑:int 的时间戳

让我们使用以毫秒为单位的纪元标准创建一个充满时间戳的假数据集。

>>>  df = sqlContext.range(1559646513000, 1559646520000)\
                    .select( (F.col('id')/1000).cast('timestamp').alias('timestamp'))
>>> df
DataFrame[timestamp: timestamp]
>>> df.show(5,False)
+-----------------------+
|timestamp              |
+-----------------------+
|2019-06-04 13:08:33    |
|2019-06-04 13:08:33.001|
|2019-06-04 13:08:33.002|
|2019-06-04 13:08:33.003|
|2019-06-04 13:08:33.004|
+-----------------------+
only showing top 5 rows

现在,转换回整数:

>>> df.select( (df.timestamp.cast('double')*1000).cast('bigint').alias('epoch') )\
      .show(5, False)
+-------------+
|epoch        |
+-------------+
|1559646513000|
|1559646513001|
|1559646513002|
|1559646513003|
|1559646513004|
+-------------+
only showing top 5 rows

【讨论】:

  • 嗨丹尼尔,谢谢你的好主意。我正在考虑将直接加入作为选项 B(虽然,正在考虑使用“外部”加入来保留所有可能的行),并查看您对重采样的评论,似乎重采样会有很多问题,所以我可能会选择加入选项。不过,您能否详细说明将时间戳转换为毫秒?怎么做?谢谢!
猜你喜欢
  • 2017-01-09
  • 1970-01-01
  • 1970-01-01
  • 2013-08-22
  • 2021-12-28
  • 1970-01-01
  • 2019-04-20
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多