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

Spark SQL数据转置求助:解决查询返回重复行问题

修正Spark SQL实现指定数据转置需求

原始数据集

branch_idbranchid_ctdiequipment_typeforecast_groupfinal_pass_yielddevice_condition
chennaichennaixd3a797.94%U
chennaichennaixd4a796.94%U
chennaichennaixd3a795.94%N

目标输出格式

branch_idbranchid_ctdiequipment_typeforecast_groupequipment_type_2final_pass_yieldfinal_pass_yield_2device_condition
chennaichennaixd3a7xd497.94%96.94%U
chennaichennaixd3a795.94%N

当前问题

现有SQL执行后出现双向重复行,多余的重复记录如下:

branch_idbranchid_ctdiequipment_typeforecast_groupequipment_type_2final_pass_yieldfinal_pass_yield_2device_condition
chennaichennaixd4a7xd396.94%97.94%U

现有SQL代码:

SELECT 
    src.branch_id, 
    src.branchid_ctdi, 
    src.equipment_type, 
    src2.equipment_type AS equipment_type_2, 
    src.forecast_group, 
    src.final_pass_yield, 
    src2.final_pass_yield AS final_pass_yield_2,
    src.device_condition
FROM 
    initial_load_stg1 src
LEFT JOIN 
    (SELECT DISTINCT branch_id, equipment_type,forecast_group,final_pass_yield,device_condition FROM initial_load_stg1 WHERE equipment_type IS NOT NULL )src2 
ON 
    src.branch_id = src2.branch_id 
    AND src.forecast_group = src2.forecast_group 
    AND src.device_condition = src2.device_condition 
    AND src.equipment_type <> src2.equipment_type
ORDER BY 
    src.branch_id, 
    src.equipment_type;

修正方案

问题根源是自连接时未限制配对的唯一性,导致每个设备类型互相配对产生双向重复。以下提供两种可靠的修正方式:

方案1:使用窗口函数标记分组内的记录顺序

通过窗口函数给每个分组(branch_id+branchid_ctdi+forecast_group+device_condition)内的记录排序,只保留分组内第一条记录作为主行,关联第二条记录作为附属字段:

WITH ranked_data AS (
    SELECT 
        branch_id,
        branchid_ctdi,
        equipment_type,
        forecast_group,
        final_pass_yield,
        device_condition,
        ROW_NUMBER() OVER (
            PARTITION BY branch_id, branchid_ctdi, forecast_group, device_condition 
            ORDER BY equipment_type
        ) AS rn
    FROM initial_load_stg1
)
SELECT 
    r1.branch_id,
    r1.branchid_ctdi,
    r1.equipment_type,
    r1.forecast_group,
    r2.equipment_type AS equipment_type_2,
    r1.final_pass_yield,
    r2.final_pass_yield AS final_pass_yield_2,
    r1.device_condition
FROM ranked_data r1
LEFT JOIN ranked_data r2
    ON r1.branch_id = r2.branch_id
    AND r1.branchid_ctdi = r2.branchid_ctdi
    AND r1.forecast_group = r2.forecast_group
    AND r1.device_condition = r2.device_condition
    AND r1.rn = 1
    AND r2.rn = 2
WHERE r1.rn = 1
ORDER BY r1.branch_id, r1.equipment_type;

方案2:修改自连接条件,强制单向配对

如果equipment_type是可比较的字符串类型,可以直接在连接条件中限制主设备类型小于关联的设备类型,避免双向配对:

SELECT 
    src.branch_id, 
    src.branchid_ctdi, 
    src.equipment_type, 
    src2.equipment_type AS equipment_type_2, 
    src.forecast_group, 
    src.final_pass_yield, 
    src2.final_pass_yield AS final_pass_yield_2,
    src.device_condition
FROM 
    initial_load_stg1 src
LEFT JOIN 
    initial_load_stg1 src2 
ON 
    src.branch_id = src2.branch_id 
    AND src.branchid_ctdi = src2.branchid_ctdi
    AND src.forecast_group = src2.forecast_group 
    AND src.device_condition = src2.device_condition 
    AND src.equipment_type < src2.equipment_type
GROUP BY 
    src.branch_id, 
    src.branchid_ctdi, 
    src.equipment_type, 
    src2.equipment_type, 
    src.forecast_group, 
    src.final_pass_yield, 
    src2.final_pass_yield,
    src.device_condition
ORDER BY 
    src.branch_id, 
    src.equipment_type;

说明

  • 方案1通用性更强,不受设备类型的排序规则限制,只要每个分组内最多两条记录,就能正确将第二条记录关联到第一条,完全避免重复。
  • 方案2代码更简洁,但依赖设备类型的字符串比较逻辑,适合明确知道设备类型可排序的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 15:37:08