【问题标题】:Mapping timeseries data to previous datapoints and averages将时间序列数据映射到以前的数据点和平均值
【发布时间】:2016-05-31 15:22:06
【问题描述】:

如果我有一个每分钟有卷的 RDD,例如

(("12:00" -> 124), ("12:01" -> 543), ("12:02" -> 102), ... )

我想将其映射到一个数据集,其中包含本分钟的音量、前一分钟的音量、前 5 分钟的平均音量。例如

(("12:00" -> (124, 300, 245.3)),
("12:01" -> (543, 124, 230.2)),
("12:02" -> (102, 543, 287.1)))

输入 RDD 可以是 RDD[(DateTime, Int)],输出可以是 RDD[(DateTime, (Int, Int, Float))]

有什么好的方法可以做到这一点?

【问题讨论】:

  • 您的数据是完整的还是可能有一些缺失的记录?
  • 可能会有差距,我会默认为零。我不介意解决方案是否能解决这个问题。
  • 在纯 scala 中,我会转换为 DateTime 并使用 SortedMap。你的数据集有多大?
  • 我想这种方法可能会有所帮助,我会在时间序列上创建窗口:stackoverflow.com/questions/23402303/…。我认为创建窗口与您尝试做的很接近,也就是说,将一个值与其系列中的上一个/下一个值组合在一起
  • 评论赞赏,有用的东西。

标签: scala apache-spark


【解决方案1】:

转换为数据框并使用窗口函数可以覆盖滞后、平均和可能的差距:

import com.github.nscala_time.time.Imports._
import org.apache.spark.sql.Row
import org.apache.spark.sql.functions.{lag, avg, when}
import org.apache.spark.sql.expressions.Window

val fmt = DateTimeFormat.forPattern("HH:mm:ss")

val rdd = sc.parallelize(Seq(
  ("12:00:00" -> 124), ("12:01:00" -> 543), ("12:02:00" -> 102),
  ("12:30:00" -> 100), ("12:31:00" -> 101)
).map{case (ds, vol) => (fmt.parseDateTime(ds), vol)})

val df = rdd
  // Convert to millis for window range
  .map{case (dt, vol) => (dt.getMillis, vol)} 
  .toDF("ts", "volume")

val w = Window.orderBy($"ts")

val transformed = df.select(
  $"ts", $"volume",
  when(
    // Check if we have data from the previous minute
    (lag($"ts", 1).over(w) - $"ts").equalTo(-60000), 
    // If so get lag otherwise 0
    lag($"volume", 1).over(w)).otherwise(0).alias("previous_volume"),
  // Average over window 
  avg($"volume").over(w.rangeBetween(-300000, 0)).alias("average"))

// Optionally go to back to RDD
transformed.map{
  case Row(ts: Long, volume: Int, previousVolume: Int, average: Double) =>
    (new DateTime(ts) -> (volume, previousVolume, average))
}

请注意,没有窗口分区的窗口函数效率很低。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-02-09
    • 1970-01-01
    • 2019-10-21
    • 2021-05-18
    • 2019-03-14
    • 2012-11-14
    • 1970-01-01
    • 2018-06-04
    相关资源
    最近更新 更多