【问题标题】:Spark: Calculate event end time on 30-minute intervals based on start time and duration values in previous rowsSpark:根据前几行中的开始时间和持续时间值以 30 分钟为间隔计算事件结束时间
【发布时间】:2019-09-05 21:46:15
【问题描述】:

我有一个带有 event_time 字段的文件,每条记录每 30 分钟生成一次,并指示事件持续了多少秒。 示例:

Event_time | event_duration_seconds
09:00      | 800
09:30      | 1800
10:00      | 2700
12:00      | 1000
13:00      | 1000

我需要将连续事件转换为仅具有持续时间的事件。输出文件应如下所示:

Event_time_start | event_time_end | event_duration_seconds
09:00            | 11:00          | 5300
12:00            | 12:30          | 1000
13:00            | 13:30          | 1000

Scala Spark 中是否有一种方法可以将数据帧记录与下一个记录进行比较?

我尝试使用foreach 循环,但不是一个好的选择,因为它需要处理大量数据

【问题讨论】:

标签: scala apache-spark dataframe hadoop apache-spark-sql


【解决方案1】:

不是一个小问题,但这里有一个解决方案,步骤如下:

  1. 使用java.time API 创建一个 UDF 以计算下一个最近的 30 分钟事件结束时间 event_ts_end
  2. 使用窗口函数lag 获取上一行的事件时间
  3. 如果与前一行的事件时间差为 30 分钟,则使用 when/otherwise 生成具有 null 值的列 event_ts_start
  4. 使用 Window 函数 last(event_ts_start, ignoreNulls=true) 用最后一个 event_ts_start 值回填 nulls
  5. event_ts_start 对数据进行分组以聚合event_durationevent_ts_end

首先,让我们组装一个示例数据集:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window
import spark.implicits._

val df = Seq(
  (101, "2019-04-01 09:00", 800),
  (101, "2019-04-01 09:30", 1800),
  (101, "2019-04-01 10:00", 2700),
  (101, "2019-04-01 12:00", 1000),
  (101, "2019-04-01 13:00", 1000),
  (220, "2019-04-02 10:00", 1500),
  (220, "2019-04-02 10:30", 2400)
).toDF("event_id", "event_time", "event_duration")

请注意,示例数据集已略微概括为包含多个事件,并使事件时间包含 date 信息以涵盖跨给定日期的事件的情况。

步骤1

import java.sql.Timestamp

def get_next_closest(seconds: Int) = udf{ (ts: Timestamp, duration: Int) =>
  import java.time.LocalDateTime
  import java.time.format.DateTimeFormatter

  val iter = Iterator.iterate(ts.toLocalDateTime)(_.plusSeconds(seconds)).
    dropWhile(_.isBefore(ts.toLocalDateTime.plusSeconds(duration)))

  iter.next.format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"))
}

步骤2 - 5

val winSpec = Window.partitionBy("event_id").orderBy("event_time")

val seconds = 30 * 60

df.
  withColumn("event_ts", to_timestamp($"event_time", "yyyy-MM-dd HH:mm")).
  withColumn("event_ts_end", get_next_closest(seconds)($"event_ts", $"event_duration")).
  withColumn("prev_event_ts", lag($"event_ts", 1).over(winSpec)).
  withColumn("event_ts_start",  when($"prev_event_ts".isNull ||
    unix_timestamp($"event_ts") - unix_timestamp($"prev_event_ts") =!= seconds, $"event_ts"
  )).
  withColumn("event_ts_start", last($"event_ts_start", ignoreNulls=true).over(winSpec)).
  groupBy($"event_id", $"event_ts_start").agg(
    sum($"event_duration").as("event_duration"), max($"event_ts_end").as("event_ts_end")
  ).show
// +--------+-------------------+--------------+-------------------+
// |event_id|     event_ts_start|event_duration|       event_ts_end|
// +--------+-------------------+--------------+-------------------+
// |     101|2019-04-01 09:00:00|          5300|2019-04-01 11:00:00|
// |     101|2019-04-01 12:00:00|          1000|2019-04-01 12:30:00|
// |     101|2019-04-01 13:00:00|          1000|2019-04-01 13:30:00|
// |     220|2019-04-02 10:00:00|          3900|2019-04-02 11:30:00|
// +--------+-------------------+--------------+-------------------+

【讨论】:

  • 感谢您的帮助,我正在尝试这种解决方案但有一些问题,我猜是因为我们的 scala 版本是 1.6
  • @mabe, to_timestamp 在 Spark 2.2 之前不可用。您可以将to_timestamp($"event_time", "yyyy-MM-dd HH:mm") 替换为from_unixtime(unix_timestamp($"event_time", "yyyy-MM-dd HH:mm"))
  • 嘿@Leo,我尝试了您的解决方案,将其调整为 spark 1.6.0 我快到了!!但是函数 last 有最后一个问题。我正在使用的 sql 库不支持 ignoreNulls 参数,因此在尝试设置 event_ts_start 值时无法过滤空值。您能否帮我选择此步骤的选项。我认为这是迄今为止我唯一缺少的东西。谢谢!
  • 嗨!最后我可以解决这个问题。我不必编写 UDF 来计算结束时间,我可以在 Hive 上下文中的 SQLText 中运行日期计算。最困难的部分是聚合记录但识别连续块。谢谢你的帮助。
  • 这里是代码`val newDF: DataFrame = records.withColumn("dt_event_start", when((lag(records("dt_real"), 1).over(winSpec).isNull) or (滞后(记录(“dt_real”),1)。over(winSpec)!==记录(“dt_prv_calc”)),记录(“dt_real”)))。 withColumn("dt_event_start", org.apache.spark.sql.functions.max("dt_event_start" ).over(winSpec))。 groupBy("dia","eutrancellfdd","dt_event_start").agg("dt_nxt_calc"->"max","re​​g" ->"sum") `
猜你喜欢
  • 2010-11-22
  • 2017-08-22
  • 2015-03-13
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-01-07
  • 2015-12-17
相关资源
最近更新 更多