MySQL通用SCD Type2存储过程数据加载异常排查求助
问题排查与修复:通用SCD Type2 MySQL存储过程
核心问题分析
你的存储过程存在致命逻辑错误,直接导致仅能处理少量记录:
- 循环对象完全错误:当前代码循环处理的是键列列表,而非源表的每一条记录。如果传入的
key_columns包含30个列,就只会处理30条(每个键列取源表第一条记录),这就是仅加载30条的根本原因。 - 键列匹配逻辑错误:SCD Type2需要用所有键列的组合匹配目标表记录,而非单独用每个键列匹配单条记录。
- 未使用拆分后的Schema名称:代码中拆分了Schema和表名,但后续动态SQL未引用Schema,跨Schema操作时会直接报错。
- SCD Type2逻辑不完整:仅标记旧记录为失效,未插入新的有效版本;且未判断记录是否真的有变化,会产生不必要的版本更新。
- 取值范围错误:仅取源表第一条记录的键值,未遍历所有源表数据。
修复后的存储过程代码
以下是修正后的通用SCD Type2存储过程,支持全量源表记录处理、多键列匹配、完整版本管理:
DELIMITER // CREATE PROCEDURE UpdateSCDType2 ( IN source_table_name VARCHAR(255), IN target_table_name VARCHAR(255), IN key_columns VARCHAR(255), -- 逗号分隔的唯一键列,如"Insurance_SK"或"Col1,Col2" IN track_columns VARCHAR(255) -- 逗号分隔的需检测变化的业务列,如"Col1,Col2,Col3" ) BEGIN DECLARE done INT DEFAULT FALSE; -- 声明游标遍历源表所有记录 DECLARE src_cursor CURSOR FOR SELECT * FROM (SELECT * FROM @full_source_table) AS src; DECLARE CONTINUE HANDLER FOR NOT FOUND SET done = TRUE; -- 拆分Schema与表名,拼接完整表名 SET @source_schema := SUBSTRING_INDEX(source_table_name, '.', 1); SET @source_table := SUBSTRING_INDEX(source_table_name, '.', -1); SET @full_source_table := CONCAT(@source_schema, '.', @source_table); SET @target_schema := SUBSTRING_INDEX(target_table_name, '.', 1); SET @target_table := SUBSTRING_INDEX(target_table_name, '.', -1); SET @full_target_table := CONCAT(@target_schema, '.', @target_table); -- 生成键列匹配条件(如`Insurance_SK` = src.`Insurance_SK`) SET @key_match := REPLACE( CONCAT(' AND ', REPLACE(key_columns, ',', ' = src.`'), ' = src.`'), ' AND ', '' ); SET @key_match := CONCAT('`', @key_match, '`'); -- 生成变化检测条件(如target.`Col1` != src.`Col1` OR ...) SET @change_check := REPLACE( CONCAT(' OR ', REPLACE(track_columns, ',', ' != src.`'), ' != src.`'), ' OR ', '' ); SET @change_check := CONCAT('target.`', @change_check, '`'); -- 遍历源表每一条记录 OPEN src_cursor; record_loop: LOOP FETCH src_cursor INTO @src_data; IF done THEN LEAVE record_loop; END IF; -- 检查目标表是否存在该键的有效版本记录 SET @exist_query := CONCAT( 'SELECT COUNT(*) INTO @record_exists FROM ', @full_target_table, ' target ', 'WHERE ', @key_match, ' AND version_ind = 1' ); PREPARE stmt FROM @exist_query; EXECUTE stmt USING @src_data; DEALLOCATE PREPARE stmt; IF @record_exists > 0 THEN -- 检查记录是否发生变化 SET @has_change_query := CONCAT( 'SELECT COUNT(*) INTO @has_change FROM ', @full_target_table, ' target ', 'WHERE ', @key_match, ' AND version_ind = 1 AND (', @change_check, ')' ); PREPARE stmt FROM @has_change_query; EXECUTE stmt USING @src_data; DEALLOCATE PREPARE stmt; IF @has_change > 0 THEN -- 将旧版本标记为失效 SET @disable_old_query := CONCAT( 'UPDATE ', @full_target_table, ' ', 'SET INSURANCE_UPDATE_DATE = CURRENT_TIMESTAMP, version_ind = 0 ', 'WHERE ', @key_match, ' AND version_ind = 1' ); PREPARE stmt FROM @disable_old_query; EXECUTE stmt USING @src_data; DEALLOCATE PREPARE stmt; -- 插入新版本记录 SET @insert_new_query := CONCAT( 'INSERT INTO ', @full_target_table, ' ', 'SELECT *, CURRENT_TIMESTAMP AS INSURANCE_CREATE_DATE, ''9999-12-31'' AS INSURANCE_UPDATE_DATE, 1 AS version_ind, 0 AS Is_Delete ', 'FROM ', @full_source_table, ' src ', 'WHERE ', @key_match ); PREPARE stmt FROM @insert_new_query; EXECUTE stmt USING @src_data; DEALLOCATE PREPARE stmt; END IF; ELSE -- 插入全新记录 SET @insert_new_query := CONCAT( 'INSERT INTO ', @full_target_table, ' ', 'SELECT *, CURRENT_TIMESTAMP AS INSURANCE_CREATE_DATE, ''9999-12-31'' AS INSURANCE_UPDATE_DATE, 1 AS version_ind, 0 AS Is_Delete ', 'FROM ', @full_source_table, ' src ', 'WHERE ', @key_match ); PREPARE stmt FROM @insert_new_query; EXECUTE stmt USING @src_data; DEALLOCATE PREPARE stmt; END IF; END LOOP; CLOSE src_cursor; -- 可选:处理源表已删除的记录(标记为失效) SET @mark_deleted_query := CONCAT( 'UPDATE ', @full_target_table, ' target ', 'SET INSURANCE_UPDATE_DATE = CURRENT_TIMESTAMP, version_ind = 0, Is_Delete = 1 ', 'WHERE version_ind = 1 ', 'AND NOT EXISTS (SELECT 1 FROM ', @full_source_table, ' src WHERE ', @key_match, ')' ); PREPARE stmt FROM @mark_deleted_query; EXECUTE stmt; DEALLOCATE PREPARE stmt; END // DELIMITER ;
关键修复点说明
- 遍历全量源表记录:使用游标遍历源表每一行数据,确保20000条记录全部被处理。
- 正确的键列匹配:将所有键列组合成匹配条件,确保唯一识别目标表中的对应记录。
- 完整SCD Type2逻辑:仅当业务字段发生变化时,标记旧版本为失效并插入新版本;无变化则不操作。
- Schema支持:动态SQL中使用完整Schema+表名,支持跨Schema操作。
- 删除记录处理:新增可选逻辑,自动标记源表已删除的目标表记录为失效。
使用示例
调用存储过程时需传入完整参数:
CALL UpdateSCDType2( 'db_source.source_insurance', 'db_target.target_insurance', 'Insurance_SK', 'Col1,Col2,Col3,Col4,Col5,Col6' );
内容的提问来源于stack exchange,提问作者Arya
相关产品推荐
相关产品推荐

