多源异构数据校验:基于Master与Driver表的SQL/C#实现问询
基于映射表和驱动表的数据校验方案
动态SQL方案(SQL Server原生实现)
前提定义(基于提供的表结构)
假设核心表结构如下:
Column_Mapping(Master表):字段名 类型 说明 Source_System VARCHAR(50) 源系统标识(如SRC_IBM) Source_Table VARCHAR(100) 源表名 Source_Column VARCHAR(100) 源列名 Target_Table VARCHAR(100) 目标表名 Target_Column VARCHAR(100) 目标列名 Is_Primary_Key BIT 是否为主键(1=是) Driver_Table:字段名 类型 说明 Selected_Source VARCHAR(50) 当前需校验的源系统
动态SQL实现代码
DECLARE @SelectedSource VARCHAR(50), @SourceTable VARCHAR(100), @TargetTable VARCHAR(100), @PrimaryKeyJoin NVARCHAR(MAX), @ColumnCompare NVARCHAR(MAX), @DiffSQL NVARCHAR(MAX) -- 1. 获取当前要校验的源系统 SELECT @SelectedSource = Selected_Source FROM Driver_Table -- 2. 获取对应源的目标表、主键连接条件、列对比语句 SELECT @SourceTable = Source_Table, @TargetTable = Target_Table, -- 构建主键连接条件(如SRC_IBM.SourceID = TRANSACTION_HEADER.TARGETID) @PrimaryKeyJoin = STRING_AGG(CONCAT(s.Source_Column, ' = t.', s.Target_Column), ' AND '), -- 构建列对比语句(如SRC_IBM.Col1 <> TRANSACTION_HEADER.TargetCol1 AS Col1_Diff) @ColumnCompare = STRING_AGG(CONCAT('CASE WHEN s.', s.Source_Column, ' <> t.', s.Target_Column, ' THEN ''Source: '' + CAST(s.', s.Source_Column, ' AS VARCHAR(MAX)) + '' | Target: '' + CAST(t.', s.Target_Column, ' AS VARCHAR(MAX)) ELSE NULL END AS ', s.Source_Column, '_Diff'), ', ') FROM Column_Mapping s WHERE s.Source_System = @SelectedSource GROUP BY s.Source_Table, s.Target_Table -- 3. 生成完整的差异查询SQL(包含源有目标无、目标有源无、列值不一致) SET @DiffSQL = N' -- 源表存在但目标表不存在的记录 SELECT ''Source_Only'' AS Diff_Type, s.*, NULL AS Target_Record, NULL AS Column_Diffs FROM ' + QUOTENAME(@SourceTable) + ' s LEFT JOIN ' + QUOTENAME(@TargetTable) + ' t ON ' + @PrimaryKeyJoin + ' WHERE t.' + (SELECT TOP 1 Target_Column FROM Column_Mapping WHERE Source_System = @SelectedSource AND Is_Primary_Key = 1) + ' IS NULL UNION ALL -- 目标表存在但源表不存在的记录 SELECT ''Target_Only'' AS Diff_Type, NULL AS Source_Record, t.*, NULL AS Column_Diffs FROM ' + QUOTENAME(@TargetTable) + ' t LEFT JOIN ' + QUOTENAME(@SourceTable) + ' s ON ' + @PrimaryKeyJoin + ' WHERE s.' + (SELECT TOP 1 Source_Column FROM Column_Mapping WHERE Source_System = @SelectedSource AND Is_Primary_Key = 1) + ' IS NULL UNION ALL -- 主键匹配但列值不一致的记录 SELECT ''Value_Mismatch'' AS Diff_Type, s.*, t.*, CONCAT('' '', STRING_AGG(CASE WHEN Diff_Col IS NOT NULL THEN Diff_Col ELSE NULL END, '', '')) AS Column_Diffs FROM ( SELECT ' + @PrimaryKeyJoin + ', ' + @ColumnCompare + ' FROM ' + QUOTENAME(@SourceTable) + ' s JOIN ' + QUOTENAME(@TargetTable) + ' t ON ' + @PrimaryKeyJoin + ' ) AS Diff_Check CROSS APPLY ( SELECT CONCAT(COLUMN_NAME, '': '', '') + COLUMN_VALUE FROM ( SELECT * FROM Diff_Check ) AS src UNPIVOT ( COLUMN_VALUE FOR COLUMN_NAME IN (' + REPLACE(@ColumnCompare, ' AS ', ' = ') + ') ) AS unpvt WHERE COLUMN_VALUE IS NOT NULL ) AS Diffs GROUP BY ' + @PrimaryKeyJoin + ' ' -- 执行动态SQL EXEC sp_executesql @DiffSQL
方案说明
- 自动适配Driver表指定的源系统,无需硬编码表名和列名
- 覆盖三类差异场景:源独有记录、目标独有记录、列值不一致
- 利用
STRING_AGG(SQL Server 2017+)简化动态语句拼接,低版本可替换为FOR XML PATH拼接
C#方案(基于ADO.NET)
核心逻辑
- 读取Driver表的当前源系统
- 从Master表拉取映射关系
- 分别查询源表和目标表数据,按主键匹配对比
- 输出结构化差异结果
示例代码
using System; using System.Collections.Generic; using System.Data.SqlClient; using System.Data; using System.Linq; public class DataValidator { private readonly string _connectionString; public DataValidator(string connectionString) { _connectionString = connectionString; } public List<ValidationDiff> ValidateCurrentSource() { var diffs = new List<ValidationDiff>(); string selectedSource = GetSelectedSource(); var mappings = GetColumnMappings(selectedSource); if (!mappings.Any()) return diffs; string sourceTable = mappings.First().SourceTable; string targetTable = mappings.First().TargetTable; var primaryKeys = mappings.Where(m => m.IsPrimaryKey).ToList(); var nonPrimaryColumns = mappings.Where(m => !m.IsPrimaryKey).ToList(); var sourcePrimaryKeys = primaryKeys.Select(m => m.SourceColumn).ToList(); var targetPrimaryKeys = primaryKeys.Select(m => m.TargetColumn).ToList(); // 获取源表和目标表数据 var sourceData = GetTableData(sourceTable, mappings.Select(m => m.SourceColumn).ToList(), sourcePrimaryKeys); var targetData = GetTableData(targetTable, mappings.Select(m => m.TargetColumn).ToList(), targetPrimaryKeys); // 检查源独有记录 foreach (var key in sourceData.Keys.Except(targetData.Keys)) { diffs.Add(new ValidationDiff { DiffType = "Source_Only", SourceRecord = sourceData[key], TargetRecord = null, ColumnDiffs = null }); } // 检查目标独有记录 foreach (var key in targetData.Keys.Except(sourceData.Keys)) { diffs.Add(new ValidationDiff { DiffType = "Target_Only", SourceRecord = null, TargetRecord = targetData[key], ColumnDiffs = null }); } // 检查列值不一致 foreach (var key in sourceData.Keys.Intersect(targetData.Keys)) { var sourceRow = sourceData[key]; var targetRow = targetData[key]; var columnMismatches = new List<string>(); foreach (var colMap in nonPrimaryColumns) { var sourceVal = sourceRow[colMap.SourceColumn]; var targetVal = targetRow[colMap.TargetColumn]; if (!Equals(sourceVal, targetVal)) { columnMismatches.Add($"{colMap.SourceColumn}: Source={sourceVal ?? "NULL"} | Target={targetVal ?? "NULL"}"); } } if (columnMismatches.Any()) { diffs.Add(new ValidationDiff { DiffType = "Value_Mismatch", SourceRecord = sourceRow, TargetRecord = targetRow, ColumnDiffs = string.Join(", ", columnMismatches) }); } } return diffs; } private string GetSelectedSource() { using (var conn = new SqlConnection(_connectionString)) { conn.Open(); var cmd = new SqlCommand("SELECT Selected_Source FROM Driver_Table", conn); return cmd.ExecuteScalar()?.ToString() ?? string.Empty; } } private List<ColumnMapping> GetColumnMappings(string sourceSystem) { var mappings = new List<ColumnMapping>(); using (var conn = new SqlConnection(_connectionString)) { conn.Open(); var cmd = new SqlCommand(@" SELECT Source_Table, Source_Column, Target_Table, Target_Column, Is_Primary_Key FROM Column_Mapping WHERE Source_System = @SourceSystem", conn); cmd.Parameters.AddWithValue("@SourceSystem", sourceSystem); using (var reader = cmd.ExecuteReader()) { while (reader.Read()) { mappings.Add(new ColumnMapping { SourceTable = reader["Source_Table"].ToString(), SourceColumn = reader["Source_Column"].ToString(), TargetTable = reader["Target_Table"].ToString(), TargetColumn = reader["Target_Column"].ToString(), IsPrimaryKey = Convert.ToBoolean(reader["Is_Primary_Key"]) }); } } } return mappings; } private Dictionary<string, Dictionary<string, object>> GetTableData(string tableName, List<string> columns, List<string> primaryKeyColumns) { var data = new Dictionary<string, Dictionary<string, object>>(); string columnList = string.Join(", ", columns.Select(c => $"[{c}]")); using (var conn = new SqlConnection(_connectionString)) { conn.Open(); var cmd = new SqlCommand($"SELECT {columnList} FROM [{tableName}]", conn); using (var reader = cmd.ExecuteReader()) { while (reader.Read()) { var rowDict = new Dictionary<string, object>(); foreach (var col in columns) { rowDict[col] = reader.IsDBNull(reader.GetOrdinal(col)) ? null : reader[col]; } // 主键拼接成字符串作为字典键 string key = string.Join("|", primaryKeyColumns.Select(col => rowDict[col] ?? "NULL")); data[key] = rowDict; } } } return data; } } // 差异结果实体类 public class ValidationDiff { public string DiffType { get; set; } public Dictionary<string, object> SourceRecord { get; set; } public Dictionary<string, object> TargetRecord { get; set; } public string ColumnDiffs { get; set; } } public class ColumnMapping { public string SourceTable { get; set; } public string SourceColumn { get; set; } public string TargetTable { get; set; } public string TargetColumn { get; set; } public bool IsPrimaryKey { get; set; } }
方案说明
- 灵活性高,可扩展添加日志、报警等功能
- 适合需要集成到现有C#数据处理流程的场景
- 可通过修改
GetTableData方法适配异构源(如Oracle、MySQL等)
替代方案
- SSIS校验组件:使用SSIS的
Lookup组件结合映射表,配置源和目标的匹配规则,输出差异记录;适合ETL流程内的校验需求。 - SQL Server CDC+映射表:启用源表和目标表的变更数据捕获,通过映射表关联CDC日志,对比增量数据的一致性;适合实时或准实时校验场景。
- 自定义校验存储过程:针对每个源系统编写专用存储过程,通过Driver表调用对应的存储过程;适合对性能要求极高的场景,避免动态SQL的开销。
内容的提问来源于stack exchange,提问作者Angel D'souza
相关产品推荐
相关产品推荐

