【问题标题】:How to merge consecutive duplicate rows in pyspark如何在pyspark中合并连续的重复行
【发布时间】:2018-05-15 05:31:08
【问题描述】:

我有一个以下格式的数据框

Col-1Col-2
a   d1
a   d2
x   d3
a   d4
f   d5
a   d6
a   d7

我想通过查看 col1 中的连续重复项来合并 col-2 中的值。我们可以看到 a 出现了两次连续重复项。它应该分别合并 d1+d2 和 d6+d7。这些列的数据类型是字符串,d1+d2 表示将字符串 d1 与 d2 连接

最终的输出应该如下图所示

Col-1Col-2
a   d1+d2
x   d3
a   d4
f   d5
a   d6+d7

【问题讨论】:

  • 你的行的数据类型是什么?它们是字符串吗?数字? +d1+d2 中是什么意思?
  • 数据类型是字符串。 + 表示简单连接 .d1+d2 是将字符串 d1 与 d2 连接
  • 您可能需要预先添加另一列来指示上一行是否重复。据我了解,在执行某些操作后,您不太可能在 RDD 中保留顺序。见here

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


【解决方案1】:

您将需要一个定义 DataFrame 顺序的列。如果一个尚不存在,您可以使用pyspark.sql.functions.monotonically_increasing_id 创建一个。

import pyspark.sql.functions as f
df = df.withColumn("id", f.monotonically_increasing_id())

接下来,您可以使用this post 中描述的技术为每组连续重复创建段:

import sys
import pyspark.sql.Window

globalWindow = Window.orderBy("id")
upToThisRowWindow = globalWindow.rowsBetween(-sys.maxsize-1, 0)

df = df.withColumn(
    "segment",
    f.sum(
        f.when(
            f.lag("Col-2", 1).over(globalWindow) != f.col("Col-2"),
            1
        ).otherwise(0)
    ).over(upToThisRowWindow)+1
)

现在您可以按段分组并使用pyspark.sql.functions.collect_list 将值收集到一个列表中并使用pyspark.sql.functions.concat() 来连接字符串:

df = df.groupBy('segment').agg(f.concat(f.collect_list('Col-2'))).drop('segment')

【讨论】:

  • 如何对连续的重复进行分组?
  • @vish 我们使用 lag 函数来检查当前行是否等于前一行。对于“Col-1”中的每个新值,都会创建一个段。
  • 这个答案有问题吗?请解释否决票。 (我可以向你保证,我不会报复投票。)如果我能理解什么是错误/不清楚的,这样我就可以修复/改进它,这将是有帮助的。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-09-02
  • 1970-01-01
  • 2018-06-16
  • 1970-01-01
  • 1970-01-01
  • 2022-01-26
  • 1970-01-01
相关资源
最近更新 更多