【问题标题】:Merge multiple spark rows inside dataframe by ID into one row based on update_time根据update_time将dataframe内的多个spark行按ID合并为一行
【发布时间】:2021-10-14 19:50:52
【问题描述】:

我们需要使用 Pyspark 将基于 ID 的多行合并为一条记录。如果该列有多个更新,那么我们必须选择对它进行最后一次更新的那个。 请注意,NULL 表示在该实例中没有对列进行更新。 因此,基本上我们必须创建一个包含对记录进行合并更新的单行。 因此,例如,如果这是数据框 ...

寻找类似的答案,但在 Pyspark .. Merge rows in a spark scala Dataframe

------------------------------------------------------------
| id       | column1          | column2         | updated_at |
------------------------------------------------------------
| 123      | update1          | <*no-update*>   | 1634228709 |   
| 123      | <*no-update*>    | 80              | 1634228724 |
| 123      | update2          | <*no-update*>   | 1634229000 |

预期输出是 -

------------------------------------------------------------
| id       | column1          | column2       | updated_at |
------------------------------------------------------------
| 123      | update2          | 80            | 1634229000 |

【问题讨论】:

  • 为什么 column2 80 在结果中。如果您想要来自最新更新的数据,它不应该为 NULL。还是有更多的规则
  • @Emerson 很抱歉造成混乱。 NULL 实际上意味着没有对该行中的列进行更新。

标签: pyspark


【解决方案1】:

假设我们的输入数据框是:

+---+-------+----+----------+
|id |col1   |col2|updated_at|
+---+-------+----+----------+
|123|null   |null|1634228709|
|123|null   |80  |1634228724|
|123|update2|90  |1634229000|
|12 |update1|null|1634221233|
|12 |null   |80  |1634228333|
|12 |update2|null|1634221220|
+---+-------+----+----------+

我们想要的是将updated_at 转换为TimestampType,然后按idupdated_at 订购desc 订购:

df = df.withColumn("updated_at", F.col("updated_at").cast(TimestampType())).orderBy(
        F.col("id"), F.col("updated_at").desc()
    )

这给了我们:

+---+-------+----+-------------------+
|id |col1   |col2|updated_at         |
+---+-------+----+-------------------+
|12 |null   |80  |2021-10-14 18:18:53|
|12 |update1|null|2021-10-14 16:20:33|
|12 |update2|null|2021-10-14 16:20:20|
|123|update2|90  |2021-10-14 18:30:00|
|123|null   |80  |2021-10-14 18:25:24|
|123|null   |null|2021-10-14 18:25:09|
+---+-------+----+-------------------+

现在在每列中获取第一个非None 值或返回None 并按id 分组:

exp = [F.first(x, ignorenulls=True).alias(x) for x in df.columns[1:]]
df = df.groupBy(F.col("id")).agg(*exp)

结果是:

+---+-------+----+-------------------+
|id |col1   |col2|updated_at         |
+---+-------+----+-------------------+
|123|update2|90  |2021-10-14 18:30:00|
|12 |update1|80  |2021-10-14 18:18:53|
+---+-------+----+-------------------+

这是完整的示例代码:

from pyspark.sql import SparkSession
import pyspark.sql.functions as F
from pyspark.sql.types import TimestampType

if __name__ == "__main__":
    spark = SparkSession.builder.master("local").appName("Test").getOrCreate()
    data = [
        (123, None, None, 1634228709),
        (123, None, 80, 1634228724),
        (123, "update2", 90, 1634229000),
        (12, "update1", None, 1634221233),
        (12, None, 80, 1634228333),
        (12, "update2", None, 1634221220),
    ]
    columns = ["id", "col1", "col2", "updated_at"]
    df = spark.createDataFrame(data, columns)
    df = df.withColumn("updated_at", F.col("updated_at").cast(TimestampType())).orderBy(
        F.col("id"), F.col("updated_at").desc()
    )
    exp = [F.first(x, ignorenulls=True).alias(x) for x in df.columns[1:]]
    df = df.groupBy(F.col("id")).agg(*exp)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-07-04
    • 2022-01-13
    • 1970-01-01
    • 2017-03-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-22
    相关资源
    最近更新 更多