【发布时间】:2020-03-25 12:28:19
【问题描述】:
对于以下数据框:
+----+--------+-------------------+----+
|user| dt| time_value|item|
+----+--------+-------------------+----+
| id1|20200101|2020-01-01 00:00:00| A|
| id1|20200101|2020-01-01 10:00:00| B|
| id1|20200101|2020-01-01 09:00:00| A|
| id1|20200101|2020-01-01 11:00:00| B|
+----+--------+-------------------+----+
我想捕获所有独特的项目,即collect_set,但保留自己的time_value
import org.apache.spark.sql.DataFrame
import org.apache.spark.sql.functions.col
import org.apache.spark.sql.functions.unix_timestamp
import org.apache.spark.sql.functions.collect_set
import org.apache.spark.sql.types.TimestampType
val timeFormat = "yyyy-MM-dd HH:mm"
val dx = Seq(("id1", "20200101", "2020-01-01 00:00", "A"), ("id1", "20200101","2020-01-01 10:00", "B"), ("id1", "20200101","2020-01-01 9:00", "A"), ("id1", "20200101","2020-01-01 11:00", "B")).toDF("user", "dt","time_value", "item").withColumn("time_value", unix_timestamp(col("time_value"), timeFormat).cast(TimestampType))
dx.show
一个
dx.groupBy("user", "dt").agg(collect_set("item")).show
+----+--------+-----------------+
|user| dt|collect_set(item)|
+----+--------+-----------------+
| id1|20200101| [B, A]|
+----+--------+-----------------+
信号从A切换到B时不保留time_value信息。如何保留item中每组的时间值信息?
是否可以在窗口函数中使用 collect_set 以达到预期的效果?目前,我只能想到:
- 使用窗口函数来确定事件对
- 过滤以更改事件
- 聚合
需要多次洗牌。或者,也可以使用 UDF (collect_list(sort_array(struct(time_value, item)))),但这似乎也很笨拙。
有没有更好的办法?
【问题讨论】:
-
您的预期结果是什么?
-
在您当前的聚合中,A 和 B 各有两个不同的“time_value”候选者,应该选择哪个?正如@Lamanus 指出的那样,很难推断出你的最终目标是什么。
标签: apache-spark apache-spark-sql user-defined-functions aggregation window-functions