You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何合并同来源的两张表?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;

代码说明

  1. full_combinations:生成所有日期与维度(Silo/GAR/Lob)的组合,确保没有缺失的记录。
  2. cumulative_recours:使用窗口函数SUM() OVER计算截至每个日期的累计recours值,实现流量转存量。
  3. latest_psap:使用LAST_VALUE(..., IGNORE NULLS)窗口函数,向前填充最新的非空PSAP值,满足缺失日期取最新存量的要求。
  4. 最后将两个结果集按维度和日期关联,得到最终期望的输出。

内容的提问来源于stack exchange,提问作者julien

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.30 05:42:03