如何合并同来源的两张表?Databricks Spark SQL技术求助
问题描述
我有一张包含流量视角(recours)和存量视角(PSAP)数据的表,拆分后已成功将流量视角字段转换为存量视角,但合并两张表时遇到困难。
尝试了两种思路:
- 生成与
columnsFluxToStock行数一致的columnsStockFinal做内连接 - 保留
columnsStock与columnsFluxToStock做全外连接
核心要求:每个缺失日期需生成记录,PSAP值取对应键(Silo/GAR/Lob)截至该日期的最新已知存量值,示例:
- 2023-03-31 / Completude / GAR1 / lob1:取2023-02-28的PSAP值2100
- 2023-03-31 / EDI / GAR1 / lob1:取2023-02-28的PSAP值1100
- 2023-04-30 / EDI / GAR1 / lob1:取2023-02-28或2023-03-31的PSAP值1100
注意:值不能硬编码。
期望结果
ViewDate Silo GAR Lob recours PSAP 2023-01-31 Completude GAR1 lob1 100 2000 2023-01-31 EDI GAR1 lob1 10 1000 2023-02-28 Completude GAR1 lob1 150 2100 2023-02-28 EDI GAR1 lob1 15 1100 2023-03-31 Completude GAR1 lob1 150 2100 2023-03-31 EDI GAR1 lob1 15 1100 2023-04-30 Completude GAR1 lob1 210 2200 2023-04-30 EDI GAR1 lob1 15 1100 2023-05-31 Completude GAR1 lob1 280 2300 2023-05-31 EDI GAR1 lob1 21 1200
现有代码(Databricks Spark SQL)
drop table if exists ColumnsMixedFluxStock; CREATE TABLE ColumnsMixedFluxStock ( ViewDate DATE, silo varchar(10), GAR varchar(10), lob varchar(10), recours int, PSAP int ); INSERT INTO ColumnsMixedFluxStock VALUES ('2023-01-31', 'EDI', 'GAR1', 'lob1', 10, 1000), ('2023-02-28', 'EDI', 'GAR1', 'lob1', 5, 1100), ('2023-05-31', 'EDI', 'GAR1', 'lob1', 6, 1200), ('2023-01-31', 'Completude', 'GAR1', 'lob1', 100, 2000), ('2023-02-28', 'Completude', 'GAR1', 'lob1', 50, 2100), ('2023-04-30', 'Completude', 'GAR1', 'lob1', 60, 2200), ('2023-05-31', 'Completude', 'GAR1', 'lob1', 70, 2300) ; with ColumnsFlux as ( select Silo, ViewDate, GAR, Lob, recours from ColumnsMixedFluxStock ), ColumnsStock as ( select Silo, ViewDate, GAR, Lob, PSAP from ColumnsMixedFluxStock ), min_max_dates AS ( SELECT MIN(ViewDate) as min_date, MAX(ViewDate) as max_date FROM ColumnsMixedFluxStock ), date_range AS ( SELECT explode(sequence(to_date(min_date), to_date(max_date), interval 1 month)) as ViewDate FROM min_max_dates ), columnsFluxToStock AS ( SELECT i1.ViewDate, i2.silo, i2.GAR, i2.Lob, SUM(coalesce(i2.recours, 0)) as recours FROM ColumnsFlux i2 cross join date_range i1 on i1.ViewDate >= i2.ViewDate GROUP BY i1.ViewDate, i2.silo, i2.GAR, i2.Lob ) /*columnsStockFinal AS ( SELECT i1.ViewDate, i2.silo, i2.GAR, i2.Lob, coalesce(last_value(i2.PSAP) OVER (PARTITION BY /*i2.ViewDate,*/ i2.silo, i2.GAR, i2.Lob ORDER BY i2.ViewDate ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW), 0) as PSAP FROM date_range i1 LEFT JOIN ColumnsStock i2 on i1.ViewDate >= i2.ViewDate ) select * from columnsStockFinal order by ViewDate, Silo, GAR, Lob*/ /*all_combinations AS ( SELECT dr.ViewDate, cs.silo, cs.GAR, cs.Lob FROM date_range dr CROSS JOIN (SELECT DISTINCT silo, GAR, Lob FROM ColumnsStock) cs ), columnsStockFinal AS ( SELECT ac.ViewDate, ac.silo, ac.GAR, ac.Lob, COALESCE(LAST_VALUE(cs.PSAP) OVER (PARTITION BY cs.silo, cs.GAR, cs.Lob ORDER BY ac.ViewDate ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW), 0) AS PSAP FROM all_combinations ac LEFT JOIN ColumnsStock cs ON ac.silo = cs.silo AND ac.GAR = cs.GAR AND ac.Lob = cs.Lob AND ac.ViewDate = cs.ViewDate ) select * from columnsStockFinal order by ViewDate, Silo, GAR, Lob*/ /* columnsStockFinal AS ( SELECT i1.ViewDate, i2.silo, i2.GAR, i2.Lob, coalesce(i2.PSAP, 0) as PSAP FROM ColumnsStock i2 cross join date_range i1 on i1.ViewDate = i2.ViewDate GROUP BY i1.ViewDate, i2.silo, i2.GAR, i2.Lob ) --select * from columnsStockFinal order by ViewDate, Silo, GAR, Lob */ /*select cs.Silo, cs.ViewDate, cs.GAR, cs.Lob, cfts.recours, cs.PSAP from columnsStockFinal cs --columnsStock cs inner join --full outer join --inner join après avoir fait une autre CTE columnsStockFinal avec un cross join de ColumnsStock et date_range columnsFluxToStock cfts on cs.Silo = cfts.Silo AND cs.ViewDate = cfts.ViewDate AND cs.GAR = cfts.GAR AND cs.Lob = cfts.Lob --where recours != 0 or PSAP != 0 order by ViewDate, Silo, GAR, Lob*/ select coalesce(cs.Silo, cfts.Silo) as Silo, coalesce(cs.ViewDate, cfts.ViewDate) ViewDate, coalesce(cs.GAR, cfts.GAR) as GAR, coalesce(cs.Lob, cfts.Lob) as Lob, cfts.recours, cs.PSAP from columnsStock cs --columnsStock cs full outer join columnsFluxToStock cfts on cs.Silo = cfts.Silo AND cs.ViewDate = cfts.ViewDate AND cs.GAR = cfts.GAR AND cs.Lob = cfts.Lob --where recours != 0 or PSAP != 0 order by ViewDate, Silo, GAR, Lob
解决方案
核心思路是先生成所有日期与维度组合的全量数据集,再分别填充recours(累计流量)和PSAP(最新存量),最后合并即可。
完整代码如下:
drop table if exists ColumnsMixedFluxStock; CREATE TABLE ColumnsMixedFluxStock ( ViewDate DATE, silo varchar(10), GAR varchar(10), lob varchar(10), recours int, PSAP int ); INSERT INTO ColumnsMixedFluxStock VALUES ('2023-01-31', 'EDI', 'GAR1', 'lob1', 10, 1000), ('2023-02-28', 'EDI', 'GAR1', 'lob1', 5, 1100), ('2023-05-31', 'EDI', 'GAR1', 'lob1', 6, 1200), ('2023-01-31', 'Completude', 'GAR1', 'lob1', 100, 2000), ('2023-02-28', 'Completude', 'GAR1', 'lob1', 50, 2100), ('2023-04-30', 'Completude', 'GAR1', 'lob1', 60, 2200), ('2023-05-31', 'Completude', 'GAR1', 'lob1', 70, 2300) ; WITH -- 获取所有维度组合 dimensions AS ( SELECT DISTINCT silo, GAR, lob FROM ColumnsMixedFluxStock ), -- 生成完整日期范围 min_max_dates AS ( SELECT MIN(ViewDate) as min_date, MAX(ViewDate) as max_date FROM ColumnsMixedFluxStock ), date_range AS ( SELECT explode(sequence(to_date(min_date), to_date(max_date), interval 1 month)) as ViewDate FROM min_max_dates ), -- 生成全量日期-维度组合 full_combinations AS ( SELECT dr.ViewDate, d.silo, d.GAR, d.lob FROM date_range dr CROSS JOIN dimensions d ), -- 计算累计recours(流量转存量) cumulative_recours AS ( SELECT fc.ViewDate, fc.silo, fc.GAR, fc.lob, SUM(cmfs.recours) OVER (PARTITION BY fc.silo, fc.GAR, fc.lob ORDER BY fc.ViewDate) AS recours FROM full_combinations fc LEFT JOIN ColumnsMixedFluxStock cmfs ON fc.silo = cmfs.silo AND fc.GAR = cmfs.GAR AND fc.lob = cmfs.lob AND cmfs.ViewDate <= fc.ViewDate ), -- 填充最新PSAP值(向前填充) latest_psap AS ( SELECT fc.ViewDate, fc.silo, fc.GAR, fc.lob, LAST_VALUE(cmfs.PSAP, IGNORE NULLS) OVER (PARTITION BY fc.silo, fc.GAR, fc.lob ORDER BY fc.ViewDate) AS PSAP FROM full_combinations fc LEFT JOIN ColumnsMixedFluxStock cmfs ON fc.silo = cmfs.silo AND fc.GAR = cmfs.GAR AND fc.lob = cmfs.lob AND fc.ViewDate = cmfs.ViewDate ) -- 合并结果 SELECT cr.ViewDate, cr.silo, cr.GAR, cr.lob, cr.recours, lp.PSAP FROM cumulative_recours cr JOIN latest_psap lp ON cr.ViewDate = lp.ViewDate AND cr.silo = lp.silo AND cr.GAR = lp.GAR AND cr.lob = lp.lob ORDER BY cr.ViewDate, cr.silo, cr.GAR, cr.lob;
代码说明
- full_combinations:生成所有日期与维度(Silo/GAR/Lob)的组合,确保没有缺失的记录。
- cumulative_recours:使用窗口函数
SUM() OVER计算截至每个日期的累计recours值,实现流量转存量。 - latest_psap:使用
LAST_VALUE(..., IGNORE NULLS)窗口函数,向前填充最新的非空PSAP值,满足缺失日期取最新存量的要求。 - 最后将两个结果集按维度和日期关联,得到最终期望的输出。
内容的提问来源于stack exchange,提问作者julien
相关产品推荐
相关产品推荐

