【问题标题】:Compare two different pyspark dataframes比较两个不同的 pyspark 数据帧
【发布时间】:2021-08-27 10:47:41
【问题描述】:

我目前正在使用需要使用 pyspark 的 API 环境。 这样,我需要在两个数据帧之间执行每日比较,以确定记录是新的、更新的和删除的。

这是两个数据框的示例:

today = spark.createDataFrame([
  [1, "Apple", 5000, "A"],
  [2, "Banana", 4000, "A"],
  [3, "Orange", 3000, "B"],
  [4, "Grape", 4500, "C"],
  [5, "Watermelon", 2000, "A"]
], ["item_id", "name", "cost", "classification"])

yesterday = spark.createDataFrame([
  [1, "Apple", 5000, "B"],
  [2, "Bananas", 4000, "A"],
  [3, "Tangerine", 3000, "B"],
  [4, "Grape", 4500, "C"]
], ["item_id", "name", "cost", "classification"])

我想比较两个数据框并确定哪些是新的,哪些是更新的。对于新项目,我很容易:

today.join(yesterday, on="item_id", how="left_anti").show()

# +---------+------------+------+----------------+
# | item_id |    name    | cost | classification |
# +---------+------------+------+----------------+
# |    5    | Watermelon | 2000 |       A        |
# +---------+------------+------+----------------+

但是对于更新的项目,我不知道如何比较这些结果。我需要为数据框的其余列获取具有不同值的所有行。 在上述情况下,我的预期结果是:

# +---------+------------+------+----------------+
# | item_id |    name    | cost | classification |
# +---------+------------+------+----------------+
# |    1    |    Apple   | 5000 |       A        |
# |    2    |   Banana   | 4000 |       A        |
# |    3    |   Orange   | 3000 |       B        |
# +---------+------------+------+----------------+

【问题讨论】:

    标签: pyspark apache-spark-sql compare


    【解决方案1】:

    使用.subtract() 方法获取today 在yesterday 中不存在的行,然后与yesterday 进行左半连接

    today.subtract(yesterday).join(yesterday, on="item_id", how="left_semi").show()
    
    # +-------+------+----+--------------+
    # |item_id|  name|cost|classification|
    # +-------+------+----+--------------+
    # |      1| Apple|5000|             A|
    # |      3|Orange|3000|             B|
    # |      2|Banana|4000|             A|
    # +-------+------+----+--------------+
    

    【讨论】:

      【解决方案2】:
      joined = today.join(yesterday, [today.item_id== yesterday.item_id] , how = 'inner' )
      

      然后应用分类不匹配的过滤器

      fitlered = joined.filter(joined.today_classification != joined.yesterday_classification)
      

      filtered 数据框是您的必需项

      【讨论】:

      • 如何确保它获取所有列不匹配的所有行?
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-05-16
      • 1970-01-01
      • 1970-01-01
      • 2019-07-21
      • 2020-06-02
      • 1970-01-01
      相关资源
      最近更新 更多