Spark SQL左连接含子查询比较条件报错,求正确实现方案
问题:基于Spark SQL实现按条件更新DataFrame的code字段
数据示例
df_a
| id | date | code |
|---|---|---|
| 1 | 2021-06-27 | A |
| 1 | 2021-12-27 | A |
| 2 | 2021-12-27 | A |
| 3 | 2022-03-21 | A |
| 3 | 2022-08-01 | A |
df_b
| id | date | code |
|---|---|---|
| 1 | 2021-05-19 | A |
| 1 | 2021-05-31 | B |
| 1 | 2021-08-27 | C |
| 3 | 2021-11-06 | X |
| 3 | 2022-02-15 | Y |
| 3 | 2022-12-30 | Z |
预期结果
| id | date | code |
|---|---|---|
| 1 | 2021-06-27 | B |
| 1 | 2021-12-27 | C |
| 2 | 2021-12-27 | A |
| 3 | 2022-03-21 | Y |
| 3 | 2022-08-01 | Y |
需求说明
用df_b中的code字段更新df_a的code字段,规则为:
- 匹配相同
id - 选取df_b中
date早于(含等于)当前df_a行date的最新记录的code - 若df_b中无符合条件的记录,保留df_a原code
原尝试SQL及报错
尝试的Spark SQL语句:
select a.id, b.code from df_a left outer join df_b on a.id = b.id and b.date = (select max(b.date) from df_b where id = a.id and date <= a.date)
报错信息:'Correlated scalar sub-queries can only be used in a Filter/Aggregate/Project and a few commands'
解决方法
Spark SQL不允许在JOIN的ON条件中使用这类关联标量子查询,以下提供两种可行方案:
方案一:窗口函数预处理+关联筛选
先对df_b按id分组并按date降序排序,再与df_a关联后,筛选出每个df_a行对应的最新df_b记录:
WITH ranked_b AS ( SELECT id, date, code, ROW_NUMBER() OVER (PARTITION BY id ORDER BY date DESC) AS rn FROM df_b ) SELECT a.id, a.date, COALESCE(b.code, a.code) AS code FROM df_a a LEFT JOIN ranked_b b ON a.id = b.id AND b.date <= a.date QUALIFY ROW_NUMBER() OVER (PARTITION BY a.id, a.date ORDER BY b.date DESC) = 1
方案二:先计算最大匹配日期再关联
先为df_a的每一行计算出df_b中符合条件的最大date,再通过该日期关联df_b获取对应code:
WITH a_max_b_date AS ( SELECT a.id, a.date AS a_date, a.code AS original_code, MAX(b.date) AS max_b_date FROM df_a a LEFT JOIN df_b b ON a.id = b.id AND b.date <= a.date GROUP BY a.id, a.date, a.code ) SELECT amb.id, amb.a_date AS date, COALESCE(b.code, amb.original_code) AS code FROM a_max_b_date amb LEFT JOIN df_b b ON amb.id = b.id AND amb.max_b_date = b.date
内容的提问来源于stack exchange,提问作者Dozel
相关产品推荐
相关产品推荐

