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

多源异构数据校验:基于Master与Driver表的SQL/C#实现问询

基于映射表和驱动表的数据校验方案

动态SQL方案(SQL Server原生实现)

前提定义(基于提供的表结构)

假设核心表结构如下:

  • Column_Mapping(Master表):
    字段名类型说明
    Source_SystemVARCHAR(50)源系统标识(如SRC_IBM)
    Source_TableVARCHAR(100)源表名
    Source_ColumnVARCHAR(100)源列名
    Target_TableVARCHAR(100)目标表名
    Target_ColumnVARCHAR(100)目标列名
    Is_Primary_KeyBIT是否为主键(1=是)
  • Driver_Table:
    字段名类型说明
    Selected_SourceVARCHAR(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)

核心逻辑

  1. 读取Driver表的当前源系统
  2. 从Master表拉取映射关系
  3. 分别查询源表和目标表数据,按主键匹配对比
  4. 输出结构化差异结果

示例代码

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等)

替代方案

  1. SSIS校验组件:使用SSIS的Lookup组件结合映射表,配置源和目标的匹配规则,输出差异记录;适合ETL流程内的校验需求。
  2. SQL Server CDC+映射表:启用源表和目标表的变更数据捕获,通过映射表关联CDC日志,对比增量数据的一致性;适合实时或准实时校验场景。
  3. 自定义校验存储过程:针对每个源系统编写专用存储过程,通过Driver表调用对应的存储过程;适合对性能要求极高的场景,避免动态SQL的开销。

内容的提问来源于stack exchange,提问作者Angel D'souza

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 13:43:13