【问题标题】:PySpark - Compare DataFramesPySpark - 比较数据帧
【发布时间】:2018-08-13 14:31:15
【问题描述】:

我是 PySpark 的新手,所以很抱歉,如果这有点简单,我发现了其他比较数据帧的问题,但不是这样的问题,因此我不认为它是重复的。 我正在尝试比较两个具有相似结构的日期框。 “名称”将是唯一的,但计数可能不同。

因此,如果计数不同,我希望它生成数据框或 python 字典。就像下面一样。关于我将如何实现这样的目标的任何想法?

DF1

+-------+---------+
|name   | count_1 |
+-------+---------+
|  Alice|   1500  |
|    Bob|   1000  |
|Charlie|   150   |
| Dexter|   100   |
+-------+---------+

DF2

+-------+---------+
|name   | count_2 |
+-------+---------+
|  Alice|   1500  |
|    Bob|   200   |
|Charlie|   150   |
| Dexter|   10    |
+-------+---------+

产生结果:

不匹配

+-------+-------------+--------------+
|name   | df1_count   | df2_count    |
+-------+-------------+--------------+
|    Bob|   1000      |    200       |
| Dexter|   100       |     10       |
+-------+-------------+--------------+

匹配

+-------+-------------+--------------+
|name   | df1_count   | df2_count    |
+-------+-------------+--------------+
|  Alice|   1500      |   1500       |
|Charlie|   150       |    150       |
+-------+-------------+--------------+

【问题讨论】:

  • 看来你应该 join 2 个数据帧和 filter 只包含 count_1 和 count_2 相等的行

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


【解决方案1】:

所以我创建了第三个 DataFrame,加入 DataFrame1 和 DataFrame2,然后通过计数字段进行过滤以检查它们是否相等:

不匹配:

df3 = df1.join(df2, [df1.name == df2.name] , how = 'inner' )
df3.filter(df3.df1_count != df3.df2_count).show()

匹配:

df3 = df1.join(df2, [df1.name == df2.name] , how = 'inner' )
df3.filter(df3.df1_count == df3.df2_count).show()

希望这对某人有用

【讨论】:

    【解决方案2】:

    对于小型 DataFrame 比较,您可以使用 chispa 库。这在测试套件中执行 DataFrame 比较时特别有用。对于大型数据集,使用连接的公认答案是最好的方法。

    在此示例中,chispa.assert_df_equality(df1, df2) 将输出以下错误消息:

    不匹配的行是红色的,匹配的行是蓝色的。 This post 有更多关于测试 PySpark 代码的信息。

    有一个很酷的库deequ 非常适合“数据单元测试”,但我不确定是否有 PySpark 实现。

    【讨论】:

      【解决方案3】:

      简单的方法是使用来自spark-extension 包的diff 转换:

      from gresearch.spark.diff import *
      
      left = spark.createDataFrame([("Alice", 1500), ("Bob", 1000), ("Charlie", 150), ("Dexter", 100)], ["name", "count"])
      right = spark.createDataFrame([("Alice", 1500), ("Bob", 200), ("Charlie", 150), ("Dexter", 10)], ["name", "count"])
      
      diff = left.diff(right, 'name')
      diff.show()
      +----+-------+----------+-----------+                                           
      |diff|   name|left_count|right_count|
      +----+-------+----------+-----------+
      |   N|  Alice|      1500|       1500|
      |   C|    Bob|      1000|        200|
      |   N|Charlie|       150|        150|
      |   C| Dexter|       100|         10|
      +----+-------+----------+-----------+
      

      这表明您在一个 DataFrame 中不匹配 (C) 和匹配 (N)。

      当然,您可以过滤以仅获取不匹配和匹配:

      diff.where(diff['diff'] == 'C').show()
      +----+------+----------+-----------+                                            
      |diff|  name|left_count|right_count|
      +----+------+----------+-----------+
      |   C|   Bob|      1000|        200|
      |   C|Dexter|       100|         10|
      +----+------+----------+-----------+
      
      diff.where(diff['diff'] == 'N').show()
      +----+-------+----------+-----------+                                           
      |diff|   name|left_count|right_count|
      +----+-------+----------+-----------+
      |   N|  Alice|      1500|       1500|
      |   N|Charlie|       150|        150|
      +----+-------+----------+-----------+
      

      虽然这是一个简单的示例,但当涉及宽模式、插入、删除和空值时,区分 DataFrame 可能会变得复杂,此解决方案完全支持。

      【讨论】:

      • 我一直在寻找资源来做这件事,我想就这个和 DataComPy 之间的比较获得意见。
      • 看起来 DataComPy 对于大多数用例来说已经足够了,但它不支持您的连接列中的空值,如您所料。我认为 DataComPy 在 diff DataFrame 中添加了一个 orderBy(join columns),当你不需要它时,这对于数十亿行来说是昂贵的。我建议为您的用例尝试它们,看看哪个结果和 API 最适合。
      【解决方案4】:

      您可以在每个数据帧之上创建一个临时视图并编写 Spark SQL 查询以进行连接。直接数据帧连接和 Spark SQL 连接是可以在这里探索的两个选项。 SQL 更直观,因此对于新手来说可能很容易。

      【讨论】:

        【解决方案5】:

        比赛

        Df1.join(Df2,Df1.col(“name”) === Df2.col(“name”) && Df1.col(“count_1”) === Df2.col(“count_2”),”inner”).drop(Df2.col(“name”)).show
        

        不匹配

        Df1.join(Df2,Df1.col(“name”) === Df2.col(“name”) && Df1.col(“count_1”) === Df2.col(“count_2”),”leftanti”).as(“Df3”).join(Df2,Df2.col(“name”) === col(“Df3.name”),”inner”).drop(Df2.col(“name”)).show
        

        【讨论】:

          猜你喜欢
          • 2018-06-08
          • 2020-01-30
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2021-05-16
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多