【问题标题】:How to join the most recent time prior to the current row time (Pyspark 2.4.4 Dataframes)如何加入当前行时间之前的最近时间(Pyspark 2.4.4 Dataframes)
【发布时间】:2020-05-28 07:29:40
【问题描述】:

我很难弄清楚如何执行以下操作:

我在 Pyspark “df1”中有 2 个数据框,如下所示:

+----+-------------+-------+
| id | SMS Created |Content|
+----+-------------+-------+
| 1  | 12:00:00    | a     |
+----+-------------+-------+
| 2  | 13:00:00    | b     |
+----+-------------+-------+
| 3  | 11:00:00    | c     |
+----+-------------+-------+

df2 看起来像这样:

+---------+----------+----+---------+
| Event   | Time     | id | Members |
+---------+----------+----+---------+
| Created | 11:30:00 | 1  | [1,2]   |
+---------+----------+----+---------+
| Updated | 11:42:00 | 1  | [1,2,3] |
+---------+----------+----+---------+
| Updated | 11:50:00 | 1  | [1,2,4] |
+---------+----------+----+---------+
| Updated | 12:50:00 | 1  | [1,2]   |
+---------+----------+----+---------+
| Created | 12:30:00 | 2  | [1,2]   |
+---------+----------+----+---------+
| Updated | 12:42:00 | 2  | [1,2,3] |
+---------+----------+----+---------+
| Updated | 12:50:00 | 2  | [1,2,4] |
+---------+----------+----+---------+
| Updated | 13:10:00 | 2  | [1,2]   |
+---------+----------+----+---------+
| Created | 10:30:00 | 3  | [1,2]   |
+---------+----------+----+---------+
| Updated | 10:42:00 | 3  | [1,2,3] |
+---------+----------+----+---------+
| Updated | 10:50:00 | 3  | [1,2,4] |
+---------+----------+----+---------+
| Updated | 12:10:00 | 2  | [1,2]   |
+---------+----------+----+---------+

df2 每次成员更改时都会更新,但消息仅发送给“创建短信”时间之前的“成员”。

请注意,在“创建短信”时间之后会有更新时间,因此在此处不使用任何类型的 MAX() 函数而无条件使用。我似乎无法理解如何做到这一点。

您将如何加入“已创建短信”之前的最新“事件”,以便表格如下所示:

+----+-------------+---------+---------+----------+---------+
| id | SMS Created | Content | Event   | Time     | Members |
+----+-------------+---------+---------+----------+---------+
| 1  | 12:00:00    | a       | Updated | 11:50:00 | [1,2.4] |
+----+-------------+---------+---------+----------+---------+
| 2  | 13:00:00    | b       | Updated | 12:50:00 | [1,2,4] |
+----+-------------+---------+---------+----------+---------+
| 3  | 11:00:00    | c       | Updated | 10:50:00 | [1,2,4] |
+----+-------------+---------+---------+----------+---------+

我正在使用带有 Dataframe API 的 Pyspark 2.4.4。任何帮助将不胜感激!

【问题讨论】:

    标签: python apache-spark pyspark apache-spark-sql


    【解决方案1】:

    welcome to SO

    试试这个:

    from pyspark.sql import functions as F
    from pyspark.sql.window import Window
    
    w=Window().partitionBy("id")
    df1.join(df2.withColumnRenamed("id","id2"), (F.col("id")==F.col("id2"))&(F.col("SMS Created")>F.col("Time"))).drop("id2")\
       .withColumn("max", F.max("Time").over(w))\
       .filter('max=Time').drop("max").orderBy("id").show()
    
    #+---+-----------+-------+-------+--------+---------+
    #| id|SMS Created|Content|  Event|    Time|  Members|
    #+---+-----------+-------+-------+--------+---------+
    #|  1|   12:00:00|      a|Updated|11:50:00|[1, 2, 4]|
    #|  2|   13:00:00|      b|Updated|12:50:00|[1, 2, 4]|
    #|  3|   11:00:00|      c|Updated|10:50:00|[1, 2, 4]|
    #+---+-----------+-------+-------+--------+---------+
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2020-05-30
      • 1970-01-01
      • 2019-02-17
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多