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

向表插入数据时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;

解决思路与方案

核心原因分析

  1. 跨数据源的sale_id冲突:单个数据源内部sale_id唯一,但加拿大和美国数据源之间可能存在相同的sale_id值,UNION ALL会直接合并这些重复值。
  2. 左连接导致行数膨胀:左连接的维度表(如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,提问作者Снежана Можейко

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 03:24:59