向表插入数据时SALE_SRC_ID重复问题的解决方案咨询
问题
从两个数据源SA_SALES_CANADA.SRC_SALES_CANADA和SA_SALES_USA.SRC_SALES_USA向目标表bl_3nf.ce_sales插入数据时,目标表的SALE_SRC_ID字段出现重复值(例如2000000000多次出现)。已知两个数据源内部的sale_id是唯一的,且已在查询中使用DISTINCT,但问题仍未解决。相关SQL代码如下:
WITH src AS ( SELECT DISTINCT SALE_ID, EVENT_DT::TIMESTAMP, CUSTOMER_ID, PRODUCT_NAME, ADDRESS, CHANNEL_ID, EMPLOYEE_ID, PAYMENT_METHOD, quantity::INT, price::NUMERIC, costs::NUMERIC, TRANSACTION_AMOUNT::NUMERIC, (price::NUMERIC - costs::NUMERIC) * quantity::INT AS REVENUE, DELIVERY_DATE::TIMESTAMP, 'SA_SALES_CANADA' AS source_system, 'SRC_CANADA' AS source_entity FROM SA_SALES_CANADA.SRC_SALES_CANADA UNION ALL SELECT DISTINCT SALE_ID, EVENT_DT::TIMESTAMP, CUSTOMER_ID, PRODUCT_NAME, ADDRESS, 'n.a.' AS CHANNEL_ID, EMPLOYEE_ID, PAYMENT_METHOD, quantity::INT, price::NUMERIC, costs::NUMERIC, TRANSACTION_AMOUNT::NUMERIC, (price::NUMERIC - costs::NUMERIC) * quantity::INT AS REVENUE, DELIVERY_DATE::TIMESTAMP, 'SA_SALES_USA' AS source_system, 'SRC_USA' AS source_entity FROM SA_SALES_USA.SRC_SALES_USA ) INSERT INTO bl_3nf.ce_sales( SALE_ID, SALE_SRC_ID, EVENT_DT, CUSTOMER_ID, PRODUCT_ID, ADDRESS_ID, CHANNEL_ID, EMPLOYEE_ID, PAYMENT_METHOD_ID, QUANTITY, PRICE, COSTS, TRANSACTION_AMOUNT, REVENUE, DELIVERY_DATE, INSERT_DT, UPDATE_DT, SOURCE_SYSTEM, SOURCE_ENTITY ) SELECT nextval('BL_3NF.SE_CE_SALES'), COALESCE(SALE_ID, 'n.a.'), COALESCE(EVENT_DT, '1900-01-01'::TIMESTAMP), CUST.CUSTOMER_ID, PROD.PRODUCT_ID, ADDR.ADDRESS_ID, CH.CHANNEL_ID, em.EMPLOYEE_ID, pm.PAYMENT_METHOD_ID, COALESCE(quantity, -1), PRICE, COSTS, TRANSACTION_AMOUNT, REVENUE, COALESCE(DELIVERY_DATE, '1900-01-01'::TIMESTAMP), NOW(), NOW(), src.source_system, src.source_entity FROM src LEFT JOIN bl_3nf.ce_addresses addr ON src.address = addr.address_src_id LEFT JOIN bl_3nf.ce_channels ch ON src.channel_id = ch.channel_src_id LEFT JOIN bl_3nf.ce_employees em ON src.employee_id = em.employee_src_id LEFT JOIN bl_3nf.ce_products prod ON src.product_NAME = prod.product_src_id LEFT JOIN bl_3nf.ce_payment_methods pm ON src.PAYMENT_METHOD = pm.payment_method_src_id LEFT JOIN BL_3NF.CE_CUSTOMERS_SDC cust ON src.customer_id = cust.customer_src_id LIMIT 50;
解决思路与方案
核心原因分析
- 跨数据源的sale_id冲突:单个数据源内部
sale_id唯一,但加拿大和美国数据源之间可能存在相同的sale_id值,UNION ALL会直接合并这些重复值。 - 左连接导致行数膨胀:左连接的维度表(如
ce_addresses、ce_channels等)中,可能存在一条源数据对应多条维度记录的情况,连接后会生成多条带相同sale_id的记录。
具体解决办法
1. 解决跨数据源的sale_id冲突
- 给不同数据源的
sale_id添加前缀,确保全局唯一:WITH src AS ( SELECT DISTINCT CONCAT('CAN_', SALE_ID) AS SALE_ID, -- 加拿大数据源前缀 EVENT_DT::TIMESTAMP, CUSTOMER_ID, PRODUCT_NAME, ADDRESS, CHANNEL_ID, EMPLOYEE_ID, PAYMENT_METHOD, quantity::INT, price::NUMERIC, costs::NUMERIC, TRANSACTION_AMOUNT::NUMERIC, (price::NUMERIC - costs::NUMERIC) * quantity::INT AS REVENUE, DELIVERY_DATE::TIMESTAMP, 'SA_SALES_CANADA' AS source_system, 'SRC_CANADA' AS source_entity FROM SA_SALES_CANADA.SRC_SALES_CANADA UNION ALL SELECT DISTINCT CONCAT('USA_', SALE_ID) AS SALE_ID, -- 美国数据源前缀 EVENT_DT::TIMESTAMP, CUSTOMER_ID, PRODUCT_NAME, ADDRESS, 'n.a.' AS CHANNEL_ID, EMPLOYEE_ID, PAYMENT_METHOD, quantity::INT, price::NUMERIC, costs::NUMERIC, TRANSACTION_AMOUNT::NUMERIC, (price::NUMERIC - costs::NUMERIC) * quantity::INT AS REVENUE, DELIVERY_DATE::TIMESTAMP, 'SA_SALES_USA' AS source_system, 'SRC_USA' AS source_entity FROM SA_SALES_USA.SRC_SALES_USA ) -- 后续插入逻辑不变 - 若允许丢弃重复项,可将
UNION ALL改为UNION,自动合并跨数据源的重复sale_id记录。
2. 消除左连接导致的行数膨胀
- 清理维度表重复数据:检查所有左连接的维度表,比如
ce_addresses中是否存在多个address_src_id等于源数据address的记录,若有则先去重(如用DELETE删除重复行,或添加唯一约束)。 - 连接后去重:在最终插入的SELECT语句外层添加
DISTINCT,确保每条sale_id只保留一条记录:INSERT INTO bl_3nf.ce_sales(...) SELECT DISTINCT nextval('BL_3NF.SE_CE_SALES'), COALESCE(SALE_ID, 'n.a.'), -- 其他字段... FROM src LEFT JOIN ... -- 原连接逻辑 - 用窗口函数筛选唯一记录:按
sale_id分组,取最新或任意一条记录:WITH joined_data AS ( SELECT nextval('BL_3NF.SE_CE_SALES') AS new_sale_id, COALESCE(SALE_ID, 'n.a.') AS sale_src_id, COALESCE(EVENT_DT, '1900-01-01'::TIMESTAMP) AS event_dt, CUST.CUSTOMER_ID, PROD.PRODUCT_ID, ADDR.ADDRESS_ID, CH.CHANNEL_ID, em.EMPLOYEE_ID, pm.PAYMENT_METHOD_ID, COALESCE(quantity, -1) AS quantity, PRICE, COSTS, TRANSACTION_AMOUNT, REVENUE, COALESCE(DELIVERY_DATE, '1900-01-01'::TIMESTAMP) AS delivery_date, NOW() AS insert_dt, NOW() AS update_dt, src.source_system, src.source_entity, ROW_NUMBER() OVER (PARTITION BY SALE_ID ORDER BY EVENT_DT DESC) AS rn FROM src LEFT JOIN bl_3nf.ce_addresses addr ON src.address = addr.address_src_id LEFT JOIN bl_3nf.ce_channels ch ON src.channel_id = ch.channel_src_id LEFT JOIN bl_3nf.ce_employees em ON src.employee_id = em.employee_src_id LEFT JOIN bl_3nf.ce_products prod ON src.product_NAME = prod.product_src_id LEFT JOIN bl_3nf.ce_payment_methods pm ON src.PAYMENT_METHOD = pm.payment_method_src_id LEFT JOIN BL_3NF.CE_CUSTOMERS_SDC cust ON src.customer_id = cust.customer_src_id ) INSERT INTO bl_3nf.ce_sales( SALE_ID, SALE_SRC_ID, EVENT_DT, CUSTOMER_ID, PRODUCT_ID, ADDRESS_ID, CHANNEL_ID, EMPLOYEE_ID, PAYMENT_METHOD_ID, QUANTITY, PRICE, COSTS, TRANSACTION_AMOUNT, REVENUE, DELIVERY_DATE, INSERT_DT, UPDATE_DT, SOURCE_SYSTEM, SOURCE_ENTITY ) SELECT new_sale_id, sale_src_id, event_dt, CUSTOMER_ID, PRODUCT_ID, ADDRESS_ID, CHANNEL_ID, EMPLOYEE_ID, PAYMENT_METHOD_ID, quantity, PRICE, COSTS, TRANSACTION_AMOUNT, REVENUE, delivery_date, insert_dt, update_dt, source_system, source_entity FROM joined_data WHERE rn = 1;
3. 从根源防止重复
给目标表bl_3nf.ce_sales的SALE_SRC_ID字段添加唯一约束,插入重复值时直接报错,便于及时发现问题:
ALTER TABLE bl_3nf.ce_sales ADD CONSTRAINT uq_ce_sales_sale_src_id UNIQUE (SALE_SRC_ID);
内容的提问来源于stack exchange,提问作者Снежана Можейко
相关产品推荐
相关产品推荐

