在这里你可以找到array functions in spark 2.4 的信息,explode_outer 是一个explode,它在一个空数组中,会产生一个'null'值的行。
想法是首先获取每个时刻,开始的消息数组,以及在每个时刻结束的消息数组(start_of 和 end_of)。
然后,我们只保留消息开始或结束的时刻,然后创建并分解以拥有一个包含 3 列的数据框,每条消息开始和结束各一列。在创建 m1 和 m2 的那一刻,将产生 2 行开始,在 m1 开始和结束的那一刻,将产生 2 行,具有 m1 星形和 m1 结束。
最后,使用窗口函数按“消息”分组并按时间排序,确保如果消息在同一时刻(同一时间)开始和结束,则首先开始。现在我们可以保证每次开始后,都会有一个结束行。
混合它们,您将获得每条消息的开头和结尾。
一个很好的思考练习。
我已经在 scala 中制作了示例,但应该很容易翻译。标记为 showAndContinue 的每一行都会在该状态下打印您的示例以显示它的作用。
val w = Window.partitionBy().orderBy("time")
val w2 = Window.partitionBy("message").orderBy($"time", desc("start_of"))
df.select($"time", $"messages", lag($"messages", 1).over(w).as("pre"), lag("messages", -1).over(w).as("post"))
.withColumn("start_of", when($"pre".isNotNull, array_except(col("messages"), col("pre"))).otherwise($"messages"))
.withColumn("end_of", when($"post".isNotNull, array_except(col("messages"), col("post"))).otherwise($"messages"))
.filter(size($"start_of") + size($"end_of") > 0)
.showAndContinue
.select(explode(array(
struct($"time", $"start_of", array().as("end_of")),
struct($"time", array().as("start_of"), $"end_of")
)).as("elem"))
.select("elem.*")
.select($"time", explode_outer($"start_of").as("start_of"), $"end_of")
.select( $"time", $"start_of", explode_outer($"end_of").as("end_of"))
.filter($"start_of".isNotNull || $"end_of".isNotNull)
.showAndContinue
.withColumn("message", when($"start_of".isNotNull, $"start_of").otherwise($"end_of"))
.showAndContinue
.select($"message", when($"start_of".isNotNull, $"time").as("starts_at"), lag($"time", -1).over(w2).as("ends_at"))
.filter($"starts_at".isNotNull)
.showAndContinue
还有桌子
+----+--------+--------+--------+--------+--------+
|time|messages| pre| post|start_of| end_of|
+----+--------+--------+--------+--------+--------+
| t01| [m1]| null|[m1, m2]| [m1]| []|
| t03|[m1, m2]| [m1]| [m2]| [m2]| [m1]|
| t04| [m2]|[m1, m2]| [m3]| []| [m2]|
| t06| [m3]| [m2]|[m3, m1]| [m3]| []|
| t07|[m3, m1]| [m3]| [m1]| [m1]| [m3]|
| t08| [m1]|[m3, m1]| [m2]| []| [m1]|
| t11| [m2]| [m1]|[m2, m4]| [m2]| []|
| t13|[m2, m4]| [m2]| [m2]| [m4]| [m4]|
| t15| [m2]|[m2, m4]| [m4]| []| [m2]|
| t20| [m4]| [m2]| []| [m4]| [m4]|
| t22|[m1, m4]| []| null|[m1, m4]|[m1, m4]|
+----+--------+--------+--------+--------+--------+
+----+--------+------+
|time|start_of|end_of|
+----+--------+------+
| t01| m1| null|
| t03| m2| null|
| t03| null| m1|
| t04| null| m2|
| t06| m3| null|
| t07| m1| null|
| t07| null| m3|
| t08| null| m1|
| t11| m2| null|
| t13| m4| null|
| t13| null| m4|
| t15| null| m2|
| t20| m4| null|
| t20| null| m4|
| t22| m1| null|
| t22| m4| null|
| t22| null| m1|
| t22| null| m4|
+----+--------+------+
+----+--------+------+-------+
|time|start_of|end_of|message|
+----+--------+------+-------+
| t01| m1| null| m1|
| t03| m2| null| m2|
| t03| null| m1| m1|
| t04| null| m2| m2|
| t06| m3| null| m3|
| t07| m1| null| m1|
| t07| null| m3| m3|
| t08| null| m1| m1|
| t11| m2| null| m2|
| t13| m4| null| m4|
| t13| null| m4| m4|
| t15| null| m2| m2|
| t20| m4| null| m4|
| t20| null| m4| m4|
| t22| m1| null| m1|
| t22| m4| null| m4|
| t22| null| m1| m1|
| t22| null| m4| m4|
+----+--------+------+-------+
+-------+---------+-------+
|message|starts_at|ends_at|
+-------+---------+-------+
| m1| t01| t03|
| m1| t07| t08|
| m1| t22| t22|
| m2| t03| t04|
| m2| t11| t15|
| m3| t06| t07|
| m4| t13| t13|
| m4| t20| t20|
| m4| t22| t22|
+-------+---------+-------+
可以优化提取在同一时刻开始和结束的所有元素,在创建的第一个表中,因此它们不必再次“匹配”开始和结束,但这取决于这是否是常见的情况,或者只是少数情况。
优化后会是这样(相同的窗口)
val dfStartEndAndFiniteLife = df.select($"time", $"messages", lag($"messages", 1).over(w).as("pre"), lag("messages", -1).over(w).as("post"))
.withColumn("start_of", when($"pre".isNotNull, array_except(col("messages"), col("pre"))).otherwise($"messages"))
.withColumn("end_of", when($"post".isNotNull, array_except(col("messages"), col("post"))).otherwise($"messages"))
.filter(size($"start_of") + size($"end_of") > 0)
.withColumn("start_end_here", array_intersect($"start_of", $"end_of"))
.withColumn("start_of", array_except($"start_of", $"start_end_here"))
.withColumn("end_of", array_except($"end_of", $"start_end_here"))
.showAndContinue
val onlyStartEndSameMoment = dfStartEndAndFiniteLife.filter(size($"start_end_here") > 0)
.select(explode($"start_end_here"), $"time".as("starts_at"), $"time".as("ends_at"))
.showAndContinue
val startEndDifferentMoment = dfStartEndAndFiniteLife
.filter(size($"start_of") + size($"end_of") > 0)
.showAndContinue
.select(explode(array(
struct($"time", $"start_of", array().as("end_of")),
struct($"time", array().as("start_of"), $"end_of")
)).as("elem"))
.select("elem.*")
.select($"time", explode_outer($"start_of").as("start_of"), $"end_of")
.select( $"time", $"start_of", explode_outer($"end_of").as("end_of"))
.filter($"start_of".isNotNull || $"end_of".isNotNull)
.showAndContinue
.withColumn("message", when($"start_of".isNotNull, $"start_of").otherwise($"end_of"))
.showAndContinue
.select($"message", when($"start_of".isNotNull, $"time").as("starts_at"), lag($"time", -1).over(w2).as("ends_at"))
.filter($"starts_at".isNotNull)
.showAndContinue
val result = onlyStartEndSameMoment.union(startEndDifferentMoment)
result.orderBy("col", "starts_at").show()
还有桌子
+----+--------+--------+--------+--------+------+--------------+
|time|messages| pre| post|start_of|end_of|start_end_here|
+----+--------+--------+--------+--------+------+--------------+
| t01| [m1]| null|[m1, m2]| [m1]| []| []|
| t03|[m1, m2]| [m1]| [m2]| [m2]| [m1]| []|
| t04| [m2]|[m1, m2]| [m3]| []| [m2]| []|
| t06| [m3]| [m2]|[m3, m1]| [m3]| []| []|
| t07|[m3, m1]| [m3]| [m1]| [m1]| [m3]| []|
| t08| [m1]|[m3, m1]| [m2]| []| [m1]| []|
| t11| [m2]| [m1]|[m2, m4]| [m2]| []| []|
| t13|[m2, m4]| [m2]| [m2]| []| []| [m4]|
| t15| [m2]|[m2, m4]| [m4]| []| [m2]| []|
| t20| [m4]| [m2]| []| []| []| [m4]|
| t22|[m1, m4]| []| null| []| []| [m1, m4]|
+----+--------+--------+--------+--------+------+--------------+
+---+---------+-------+
|col|starts_at|ends_at|
+---+---------+-------+
| m4| t13| t13|
| m4| t20| t20|
| m1| t22| t22|
| m4| t22| t22|
+---+---------+-------+
+----+--------+--------+--------+--------+------+--------------+
|time|messages| pre| post|start_of|end_of|start_end_here|
+----+--------+--------+--------+--------+------+--------------+
| t01| [m1]| null|[m1, m2]| [m1]| []| []|
| t03|[m1, m2]| [m1]| [m2]| [m2]| [m1]| []|
| t04| [m2]|[m1, m2]| [m3]| []| [m2]| []|
| t06| [m3]| [m2]|[m3, m1]| [m3]| []| []|
| t07|[m3, m1]| [m3]| [m1]| [m1]| [m3]| []|
| t08| [m1]|[m3, m1]| [m2]| []| [m1]| []|
| t11| [m2]| [m1]|[m2, m4]| [m2]| []| []|
| t15| [m2]|[m2, m4]| [m4]| []| [m2]| []|
+----+--------+--------+--------+--------+------+--------------+
+----+--------+------+
|time|start_of|end_of|
+----+--------+------+
| t01| m1| null|
| t03| m2| null|
| t03| null| m1|
| t04| null| m2|
| t06| m3| null|
| t07| m1| null|
| t07| null| m3|
| t08| null| m1|
| t11| m2| null|
| t15| null| m2|
+----+--------+------+
+----+--------+------+-------+
|time|start_of|end_of|message|
+----+--------+------+-------+
| t01| m1| null| m1|
| t03| m2| null| m2|
| t03| null| m1| m1|
| t04| null| m2| m2|
| t06| m3| null| m3|
| t07| m1| null| m1|
| t07| null| m3| m3|
| t08| null| m1| m1|
| t11| m2| null| m2|
| t15| null| m2| m2|
+----+--------+------+-------+
+-------+---------+-------+
|message|starts_at|ends_at|
+-------+---------+-------+
| m1| t01| t03|
| m1| t07| t08|
| m2| t03| t04|
| m2| t11| t15|
| m3| t06| t07|
+-------+---------+-------+
+---+---------+-------+
|col|starts_at|ends_at|
+---+---------+-------+
| m1| t01| t03|
| m1| t07| t08|
| m1| t22| t22|
| m2| t03| t04|
| m2| t11| t15|
| m3| t06| t07|
| m4| t13| t13|
| m4| t20| t20|
| m4| t22| t22|
+---+---------+-------+