【问题标题】:Spark SQL - Keep one result after joinSpark SQL - 加入后保留一个结果
【发布时间】:2023-03-15 21:45:01
【问题描述】:

我有两个数据框,我试图将它们连接在一起,但我意识到在我最初的实现中,我得到了不希望的结果:

// plain_txns_df.show(false)
+------------+---------------+---------------+-----------+
|txn_date_at |account_number |merchant_name  |txn_amount |
+------------+---------------+---------------+-----------+
|2020-04-08  |1234567        |Starbucks      |2.02       |
|2020-04-14  |1234567        |Starbucks      |2.86       |
|2020-04-14  |1234567        |Subway         |12.02      |
|2020-04-14  |1234567        |Amazon         |3.21       |
+------------+---------------+---------------+-----------+
// richer_txns_df.show(false)
+----------+-------+----------------------+-------------+
|TXN_DT    |ACCT_NO|merch_name            |merchant_city|
+----------+-------+----------------------+-------------+
|2020-04-08|1234567|Subway                |Toronto      |
|2020-04-14|1234567|Subway                |Toronto      |
+----------+-------+----------------------+-------------+

从以上两个数据框中,我的目标是丰富与商家城市的普通交易,对于 7 天窗口内的交易(即来自更丰富的交易数据框的交易日期应该在普通日期和普通日期之间)日期 - 7 天。

一开始我认为这很简单,就这样加入了数据(我知道范围加入):

spark.sql(
    """
      | SELECT
      | plain.txn_date_at,
      | plain.account_number,
      | plain.merchant_name,
      | plain.txn_amount,
      | richer.merchant_city
      | FROM plain_txns_df plain
      | LEFT JOIN richer_txns_df richer
      | ON plain.account_number = richer.ACCT_NO
      | AND plain.merchant_name = richer.merch_name
      | AND richer.txn_date BETWEEN date_sub(plain.txn_date_at, 7) AND plain.txn_date_at
    """.stripMargin)

但是,当使用上述方法时,我得到了 4 月 14 日交易的重复结果,因为商家详细信息和帐户详细信息与 8 日的更丰富的记录相匹配,并且符合日期范围:

+------------+---------------+---------------+-----------+-------------+
|txn_date_at |account_number |merchant_name  |txn_amount |merchant_city|
+------------+---------------+---------------+-----------+-------------+
|2020-04-08  |1234567        |Starbucks      |2.02       |Toronto      |
|2020-04-14  |1234567        |Starbucks      |2.86       |Toronto      | // Apr-08 Richer record
|2020-04-14  |1234567        |Starbucks      |2.86       |Toronto      |
+------------+---------------+---------------+-----------+-------------+

有没有一种方法可以为我的普通 DataFrame 中的每个值获取 一个 记录(即在上述结果集中为第 14 个获取一个记录)? 我尝试在加入后运行一个不同的,这解决了这个问题,但我意识到如果同一天有两笔同一商家的交易,我会失去这些。

我正在考虑将更丰富的表移动到子查询中,然后在其中应用日期过滤器,但我不知道如何将事务日期过滤器值传递到此查询中:(。类似于以下内容,但它无法识别普通交易日期:

spark.sql(
    """
      | SELECT
      | plain.txn_date_at,
      | plain.account_number,
      | plain.merchant_name,
      | plain.txn_amount,
      | richer2.merchant_city
      | FROM plain_txns_df plain
      | LEFT JOIN ( 
      |    SELECT ACCT_NO, merch_name from richer_txns_df
      |    WHERE txn_date BETWEEN date_sub(plain.txn_date_at, 7) AND plain.txn_date_at
      | ) richer2
      | ON plain.account_number = richer2.ACCT_NO
      | AND plain.merchant_name = richer2.merch_name
    """.stripMargin)

【问题讨论】:

  • 我认为您还需要在加入条件中添加日期。所以它将txn_date_at、account_number和merchant_name作为一个复合键,你不会得到重复

标签: apache-spark apache-spark-sql correlated-subquery


【解决方案1】:

我认为需要做的第一件事是在plain_txns_df 上创建一个唯一键,这使得在尝试聚合/比较它们时可以将它们彼此区分开来。

import org.apache.spark.sql.functions._
plainDf.withColumn("id", monotonically_increasing_id())

这样您就可以继续执行您发布的第一个查询(加上id 列),该查询返回重复项:

spark.sql("""
    SELECT
    plain.id,
    plain.txn_date_at,
    plain.account_number,
    plain.merchant_name,
    plain.txn_amount,
    richer.merchant_city,
    richer.txn_dt
    FROM plain_txns_df plain
    INNER JOIN richer_txns_df richer
    ON plain.account_number = richer.acc_no
    AND plain.merchant_name = richer.merch_name
    AND richer.txn_dt BETWEEN date_sub(plain.txn_date_at, 7) AND plain.txn_date_at
  """.stripMargin).createOrReplaceTempView("foo")

接下来是通过获取给定id 的最新richer_txns_df.txn_dt 日期记录来对上述数据帧进行重复数据删除。

spark.sql("""
    SELECT
    f1.txn_date_at,
    f1.account_number,
    f1.merchant_name,
    f1.txn_amount,
    f1.merchant_city
    FROM foo f1
    LEFT JOIN foo f2
    ON f2.id = f1.id
    AND f2.txn_dt > f1.txn_dt
    WHERE f2.id IS NULL
  """.stripMargin).show

【讨论】:

  • 有道理!您认为这可以使用窗口排名函数而不是将表格连接到自身来完成吗?
猜你喜欢
  • 2013-01-25
  • 2019-03-03
  • 2021-03-02
  • 1970-01-01
  • 1970-01-01
  • 2021-12-19
  • 1970-01-01
  • 2021-01-21
  • 1970-01-01
相关资源
最近更新 更多