Spark SQL数据转置求助:解决查询返回重复行问题
修正Spark SQL实现指定数据转置需求
原始数据集
| branch_id | branchid_ctdi | equipment_type | forecast_group | final_pass_yield | device_condition |
|---|---|---|---|---|---|
| chennai | chennai | xd3 | a7 | 97.94% | U |
| chennai | chennai | xd4 | a7 | 96.94% | U |
| chennai | chennai | xd3 | a7 | 95.94% | N |
目标输出格式
| branch_id | branchid_ctdi | equipment_type | forecast_group | equipment_type_2 | final_pass_yield | final_pass_yield_2 | device_condition |
|---|---|---|---|---|---|---|---|
| chennai | chennai | xd3 | a7 | xd4 | 97.94% | 96.94% | U |
| chennai | chennai | xd3 | a7 | 95.94% | N |
当前问题
现有SQL执行后出现双向重复行,多余的重复记录如下:
| branch_id | branchid_ctdi | equipment_type | forecast_group | equipment_type_2 | final_pass_yield | final_pass_yield_2 | device_condition |
|---|---|---|---|---|---|---|---|
| chennai | chennai | xd4 | a7 | xd3 | 96.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
相关产品推荐
相关产品推荐

