【问题标题】:Compare two dataframes Pyspark比较两个数据框 Pyspark
【发布时间】:2020-06-02 08:55:24
【问题描述】:

我正在尝试比较两个具有相同列数的数据框,即在两个数据框中以 id 作为关键列的 4 列

df1 = spark.read.csv("/path/to/data1.csv")
df2 = spark.read.csv("/path/to/data2.csv")

现在我想将新列附加到 DF2,即 column_names,它是与 df1 具有不同值的列的列表

df2.withColumn("column_names",udf())

DF1

+------+---------+--------+------+
|   id | |name  | sal  | Address |
+------+---------+--------+------+
|     1|  ABC   | 5000 | US      |
|     2|  DEF   | 4000 | UK      |
|     3|  GHI   | 3000 | JPN     |
|     4|  JKL   | 4500 | CHN     |
+------+---------+--------+------+

DF2:

+------+---------+--------+------+
|   id | |name  | sal  | Address |
+------+---------+--------+------+
|     1|  ABC   | 5000 | US      |
|     2|  DEF   | 4000 | CAN     |
|     3|  GHI   | 3500 | JPN     |
|     4|  JKL_M | 4800 | CHN     |
+------+---------+--------+------+

现在我想要 DF3

DF3:

+------+---------+--------+------+--------------+
|   id | |name  | sal  | Address | column_names |
+------+---------+--------+------+--------------+
|     1|  ABC   | 5000 | US      |  []          |
|     2|  DEF   | 4000 | CAN     |  [address]   |
|     3|  GHI   | 3500 | JPN     |  [sal]       |
|     4|  JKL_M | 4800 | CHN     |  [name,sal]  |
+------+---------+--------+------+--------------+

我看到了这个 SO 问题,How to compare two dataframe and print columns that are different in scala。试过了,结果不一样。

我正在考虑通过将每个数据帧中的行传递给 udf 并逐列比较并返回列列表来使用 UDF 函数。但是,为此,两个数据帧都应该按排序顺序排列,以便将相同的 id 行发送到 udf。排序在这里是昂贵的操作。有什么解决办法吗?

【问题讨论】:

  • 你想要 pyspark 还是 spark 的解决方案?是指scala还是python?
  • 我在 python 中寻找解决方案
  • Sorting is costly operation here. - 在这种情况下,我认为没有比排序更好的方法了
  • 这里不需要UDF。在id 上使用左连接,然后比较列值并创建新列column_names。

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


【解决方案1】:

假设我们可以使用 id 来加入这两个数据集,我认为不需要 UDF。这可以通过使用内连接、array 和 array_remove 等函数来解决。

首先让我们创建两个数据集:

df1 = spark.createDataFrame([
  [1, "ABC", 5000, "US"],
  [2, "DEF", 4000, "UK"],
  [3, "GHI", 3000, "JPN"],
  [4, "JKL", 4500, "CHN"]
], ["id", "name", "sal", "Address"])

df2 = spark.createDataFrame([
  [1, "ABC", 5000, "US"],
  [2, "DEF", 4000, "CAN"],
  [3, "GHI", 3500, "JPN"],
  [4, "JKL_M", 4800, "CHN"]
], ["id", "name", "sal", "Address"])

首先我们在两个数据集之间进行内部连接,然后为除id 之外的每一列生成条件df1[col] != df2[col]。当列不相等时,我们返回列名,否则返回空字符串。条件列表将包含一个数组的项,最后我们从中删除空项:

from pyspark.sql.functions import col, array, when, array_remove

# get conditions for all columns except id
conditions_ = [when(df1[c]!=df2[c], lit(c)).otherwise("") for c in df1.columns if c != 'id']

select_expr =[
                col("id"), 
                *[df2[c] for c in df2.columns if c != 'id'], 
                array_remove(array(*conditions_), "").alias("column_names")
]

df1.join(df2, "id").select(*select_expr).show()

# +---+-----+----+-------+------------+
# | id| name| sal|Address|column_names|
# +---+-----+----+-------+------------+
# |  1|  ABC|5000|     US|          []|
# |  3|  GHI|3500|    JPN|       [sal]|
# |  2|  DEF|4000|    CAN|   [Address]|
# |  4|JKL_M|4800|    CHN| [name, sal]|
# +---+-----+----+-------+------------+

【讨论】:

  • 这里不需要otherwise("")。将它们保留为空值,然后使用 filter 或 array_except 删除空值。
  • array_except 只能与 array_except(array(*conditions_), array(lit(None))) 一起使用,这会引入额外的开销来创建一个新数组而不需要它。至于filter,我认为 pyspark 只能通过expr 或selectExpr 获得,或者至少databricks 拒绝将它包含在from pyspark.sql.functions import filter 中,并且确实似乎没有出现在functions
【解决方案2】:

Python:我之前的 scala 代码的 PySpark 版本。

import pyspark.sql.functions as f

df1 = spark.read.option("header", "true").csv("test1.csv")
df2 = spark.read.option("header", "true").csv("test2.csv")

columns = df1.columns
df3 = df1.alias("d1").join(df2.alias("d2"), f.col("d1.id") == f.col("d2.id"), "left")

for name in columns:
    df3 = df3.withColumn(name + "_temp", f.when(f.col("d1." + name) != f.col("d2." + name), f.lit(name)))


df3.withColumn("column_names", f.concat_ws(",", *map(lambda name: f.col(name + "_temp"), columns))).select("d1.*", "column_names").show()

Scala:这是我解决问题的最佳方法。

val df1 = spark.read.option("header", "true").csv("test1.csv")
val df2 = spark.read.option("header", "true").csv("test2.csv")

val columns = df1.columns
val df3 = df1.alias("d1").join(df2.alias("d2"), col("d1.id") === col("d2.id"), "left")

columns.foldLeft(df3) {(df, name) => df.withColumn(name + "_temp", when(col("d1." + name) =!= col("d2." + name), lit(name)))}
  .withColumn("column_names", concat_ws(",", columns.map(name => col(name + "_temp")): _*))
  .show(false)

首先,我将两个数据框加入df3 并使用来自df1 的列。当df1 和df2 具有相同的id 和其他列值时,通过将具有列名称值的临时列向左折叠到df3。

在那之后,那些列名的concat_ws 和空值都消失了,只剩下列名。

+---+----+----+-------+------------+
|id |name|sal |Address|column_names|
+---+----+----+-------+------------+
|1  |ABC |5000|US     |            |
|2  |DEF |4000|UK     |Address     |
|3  |GHI |3000|JPN    |sal         |
|4  |JKL |4500|CHN    |name,sal    |
+---+----+----+-------+------------+

与您的预期结果唯一不同的是输出不是列表而是字符串。

附言我忘记使用 PySpark 但这是正常的火花,抱歉。

【讨论】:

    【解决方案3】:

    这是您使用 UDF 的解决方案,我已经动态更改了第一个 dataframe 名称,以便在检查时不会有歧义。浏览下面的代码,如果有任何问题,请告诉我。

    >>> from pyspark.sql.functions import *
    >>> df.show()
    +---+----+----+-------+
    | id|name| sal|Address|
    +---+----+----+-------+
    |  1| ABC|5000|     US|
    |  2| DEF|4000|     UK|
    |  3| GHI|3000|    JPN|
    |  4| JKL|4500|    CHN|
    +---+----+----+-------+
    
    >>> df1.show()
    +---+----+----+-------+
    | id|name| sal|Address|
    +---+----+----+-------+
    |  1| ABC|5000|     US|
    |  2| DEF|4000|    CAN|
    |  3| GHI|3500|    JPN|
    |  4|JKLM|4800|    CHN|
    +---+----+----+-------+
    
    >>> df2 = df.select([col(c).alias("x_"+c) for c in df.columns])
    >>> df3 = df1.join(df2, col("id") == col("x_id"), "left")
    
     //udf declaration 
    
    >>> def CheckMatch(Column,r):
    ...     check=''
    ...     ColList=Column.split(",")
    ...     for cc in ColList:
    ...             if(r[cc] != r["x_" + cc]):
    ...                     check=check + "," + cc
    ...     return check.replace(',','',1).split(",")
    
    >>> CheckMatchUDF = udf(CheckMatch)
    
    //final column that required to select
    >>> finalCol = df1.columns
    >>> finalCol.insert(len(finalCol), "column_names")
    
    >>> df3.withColumn("column_names", CheckMatchUDF(lit(','.join(df1.columns)),struct([df3[x] for x in df3.columns])))
           .select(finalCol)
           .show()
    +---+----+----+-------+------------+
    | id|name| sal|Address|column_names|
    +---+----+----+-------+------------+
    |  1| ABC|5000|     US|          []|
    |  2| DEF|4000|    CAN|   [Address]|
    |  3| GHI|3500|    JPN|       [sal]|
    |  4|JKLM|4800|    CHN| [name, sal]|
    +---+----+----+-------+------------+
    

    【讨论】:

    【解决方案4】:

    您可以通过 spark-extension 包在 PySpark 和 Scala 中为您构建查询。 它提供的diff 转换正是这样做的。

    from gresearch.spark.diff import *
    
    options = DiffOptions().with_change_column('changes')
    df1.diff_with_options(df2, options, 'id').show()
    +----+-----------+---+---------+----------+--------+---------+------------+-------------+
    |diff|    changes| id|left_name|right_name|left_sal|right_sal|left_Address|right_Address|
    +----+-----------+---+---------+----------+--------+---------+------------+-------------+
    |   N|         []|  1|      ABC|       ABC|    5000|     5000|          US|           US|
    |   C|  [Address]|  2|      DEF|       DEF|    4000|     4000|          UK|          CAN|
    |   C|      [sal]|  3|      GHI|       GHI|    3000|     3500|         JPN|          JPN|
    |   C|[name, sal]|  4|      JKL|     JKL_M|    4500|     4800|         CHN|          CHN|
    +----+-----------+---+---------+----------+--------+---------+------------+-------------+
    
    

    虽然这是一个简单的示例,但当涉及到宽模式、插入、删除和空值时,区分 DataFrame 可能会变得复杂。该软件包已经过充分测试,因此您不必担心自己会正确查询。

    【讨论】:

      【解决方案5】:

      了解它在 pyspark 上的工作方式

      从 gresearch.spark.diff 导入 *

      options = DiffOptions().with_change_column('changes') df1.diff_with_options(df2, options, 'id').show() +----+------------+---+----------+----------+------- -+---------+------------+--------------+ |差异|变化| id|left_name|right_name|left_sal|right_sal|left_Address|right_Address| +----+------------+---+----------+----------+------- -+---------+------------+--------------+ | N| []| 1| ABC| ABC| 5000| 5000|美国|美国| | C| [地址]| 2|防御|防御| 4000| 4000|英国|罐头| | C| [萨尔]| 3|全球健康指数|全球健康指数| 3000| 3500|日本|日本| | C|[姓名,萨尔]| 4| JKL| JKL_M| 4500| 4800|中国|中国| +----+------------+---+----------+----------+------- -+---------+------------+-------------+

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2017-07-30
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2021-05-16
        相关资源
        最近更新 更多