【问题标题】:grouping pyspark rows based on condtion根据条件对 pyspark 行进行分组
【发布时间】:2021-06-14 22:12:04
【问题描述】:

我有这个有 6 列的表,我想根据“记录”字段按“ID1”和“ID2”对行进行分组。我的记录字段是“IN”或“OUT”,它们按日期排序。

这是我的输入样本...

data = [("ACC.PXP","7246","2020-02-24T14:49:00",None,None,'IN'),
    ("ACC.PXP","7246","2021-03-09T08:20:00","Hospital","Foundation","OUT"),
    ("ACC.PXP","7246","2021-04-05T17:17:00","Hospital","Foundation","IN")
       ] 
df = spark.createDataFrame(data=data,schema=['ID1','ID2','date','type','name','record'])
df.show(truncate=False)

+-------+----+-------------------+--------+----------+------+
|ID1    |ID2 |date               |type    |name      |record|
+-------+----+-------------------+--------+----------+------+
|ACC.PXP|7246|2020-02-24T14:49:00|null    |null      |IN    |
|ACC.PXP|7246|2021-03-09T08:20:00|Hospital|Foundation|OUT   |
|ACC.PXP|7246|2021-04-05T17:17:00|Hospital|Foundation|IN    |

这就是我想要的结果

data2 = [("ACC.PXP","7246","2020-02-24T14:49:00",None,None, "2021-03-09T08:20:00","Hospital","Foundation"),
    ("ACC.PXP","7246","2021-04-05T17:17:00","Hospital","Foundation", None,None,None)
  ]
 
df2 = spark.createDataFrame(data=data2,schema=['ID1','ID2','date','type','name','date1','type1','name1'])
df2.show(truncate=False)

+-------+----+-------------------+--------+----------+-------------------+--------+----------+
|ID1    |ID2 |date               |type    |name      |date1              |type1   |name1     |
+-------+----+-------------------+--------+----------+-------------------+--------+----------+
|ACC.PXP|7246|2020-02-24T14:49:00|null    |null      |2021-03-09T08:20:00|Hospital|Foundation|
|ACC.PXP|7246|2021-04-05T17:17:00|Hospital|Foundation|null               |null    |null      |
+-------+----+-------------------+--------+----------+-------------------+--------+----------+

【问题讨论】:

  • @sammywemmy 你知道如何解决这个问题吗?
  • 嗨@ScootCork 我在stackoverflow.com/questions/57435858/… 看到了您的回答,我的问题与您的回答相似。你觉得你能帮上忙吗?谢谢

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


【解决方案1】:

您可以按 id 和“计数”分组,并按如下方式透视“记录”列:

import pyspark.sql.functions as F
from pyspark.sql import Window
w = Window.partitionBy("ID1", "ID2", "record").orderBy("date")
df1 = (df
   .withColumn("count", F.row_number().over(w))
   .groupBy("ID1", "ID2", "count")
   .pivot("record")
   .agg(F.first("date"), F.first( "type"), F.first("name"))
   .select("ID1", "ID2", 
           F.col("IN_first(date)").alias("date"),  
           F.col("IN_first(type)").alias("type"), 
           F.col("IN_first(name)").alias("name"),
           F.col("OUT_first(date)").alias("date1"),
           F.col("OUT_first(date)").alias("type1"),
           F.col("OUT_first(name)").alias("name1"))
  )

这将生成所需的输出表。但是,只是一个与您的解决方案相同的警告,此解决方案假定对于每个 ID 对,第一个有日期的条目用于记录 = IN,并且如果按日期排序,则行遵循 IN-OUT-IN-OUT... 序列.否则,此解决方案将无法正常工作。

【讨论】:

  • 谢谢@安娜
  • 嗨@Anna,您的代码在按日期排序时适用于 IN-ONT-IN... 序列。但是,如果我有坏数据并且第一条记录从 OUT 开始,我想做和以前一样的事情怎么办。说我的输入是:``` ```
【解决方案2】:

我相信有人会想出一个更快、更短、更优雅的 pyspark 代码。但是这个也可以。

from pyspark.sql import functions as F
from pyspark.sql import Window
data = [("ACC.PXP","7246","2020-02-24T14:49:00",None,None,'IN'),
    ("ACC.PXP","7246","2021-03-09T08:20:00","Hospital","Foundation","OUT"),
    ("ACC.PXP","7246","2021-04-05T17:17:00","Hospital","Foundation","IN")
       ] 
sdf = spark.createDataFrame(data=data,schema=['ID1','ID2','date','type','name','record'])
sdf.show(truncate=False)

## split the data frame into two based on record type and give it a one 
sdf_1 = sdf.filter("record == 'IN'").withColumn('ones', F.lit(1))
sdf_2 = (sdf.filter("record == 'OUT'").withColumnRenamed('date', 'date1')\
                                      .withColumnRenamed('type', 'type1')\
                                        .withColumnRenamed('name', 'name1')\
                                          .withColumn('ones', F.lit(1))
)

## partition it by id1 and id2 
windowSpec1 = Window.partitionBy("ID1","ID2").orderBy("date")
windowSpec2 = Window.partitionBy("ID1","ID2").orderBy("date1")

## creat a count column to count the number of 'IN' and 'OUT'
sdf_1 = (sdf_1.withColumn('counter', F.sum('ones').over(windowSpec1))\
              .drop("ones", "record")
        )
sdf_2 = (sdf_2.withColumn('counter', F.sum('ones').over(windowSpec2))\
               .withColumnRenamed('ID1','r_ID1')\
                 .withColumnRenamed('ID2','r_ID2')\
                   .withColumnRenamed('counter','r_counter')\
                     .drop("ones", "record")
         )
## merge the two dataframes back
sdf_merged = (sdf_1.join(sdf_2, ( sdf_1.ID1 == sdf_2.r_ID1) & 
                                        (sdf_1.ID2 == sdf_2.r_ID2)  & 
                                        (sdf_1.counter  == sdf_2.r_counter), 
                                        how ='left')\
                                .drop(sdf_2.r_ID1).drop(sdf_2.r_ID2).drop(sdf_2.r_counter).drop(sdf_1.counter)\
                                .orderBy(F.asc('date'))
             )
sdf_merged.show()

+-------+----+-------------------+--------+----------+-------------------+--------+----------+
|    ID1| ID2|               date|    type|      name|              date1|   type1|     name1|
+-------+----+-------------------+--------+----------+-------------------+--------+----------+
|ACC.PXP|7246|2020-02-24T14:49:00|    null|      null|2021-03-09T08:20:00|Hospital|Foundation|
|ACC.PXP|7246|2021-04-05T17:17:00|Hospital|Foundation|               null|    null|      null|
+-------+----+-------------------+--------+----------+-------------------+--------+----------+

【讨论】:

    猜你喜欢
    • 2021-08-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-03-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多