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

基于C# Dapper的单事务批量增删改动态SQL实现方案

基于Dapper的单事务批量同步解决方案

针对你的需求,核心解决方案是临时表+SqlBulkCopy+批量SQL操作,既满足单事务要求,又能规避参数数量限制,同时通过动态生成安全的SQL实现防注入,完全不需要依赖EF。

核心思路

  1. 利用Session级临时表批量导入待同步数据,替代逐个生成参数的方式,避开SQL Server参数数量上限(默认2100个)。
  2. 通过反射动态生成临时表结构、目标表的增删改SQL,保证DTO与数据库表的一一对应。
  3. 所有操作在单个数据库事务内执行,确保同步的原子性。
  4. 采用SqlBulkCopy高效导入数据,配合批量JOIN语句完成增删改,大幅提升处理效率。

关键实现步骤

1. 类型映射准备

定义.NET类型到SQL Server类型的映射字典,用于动态生成临时表结构:

private static readonly Dictionary<Type, string> SqlTypeMappings = new()
{
    { typeof(int), "INT" },
    { typeof(int?), "INT NULL" },
    { typeof(string), "NVARCHAR(MAX)" },
    { typeof(DateTime), "DATETIME2" },
    { typeof(DateTime?), "DATETIME2 NULL" },
    { typeof(bool), "BIT" },
    { typeof(bool?), "BIT NULL" },
    { typeof(decimal), "DECIMAL(18,2)" },
    { typeof(decimal?), "DECIMAL(18,2) NULL" },
    // 按需补充项目中用到的其他类型
};

2. 动态生成临时表SQL

通过反射获取DTO的属性信息,生成对应的临时表创建语句:

// 生成包含所有字段的临时表(用于新增、更新)
private string GenerateCreateTempTableSql(Type dtoType)
{
    var properties = dtoType.GetProperties(BindingFlags.Public | BindingFlags.Instance);
    var columns = string.Join(",\n", properties.Select(p => 
        $"[{p.Name}] {SqlTypeMappings[Nullable.GetUnderlyingType(p.PropertyType) ?? p.PropertyType]}"));
    return $"CREATE TABLE #{dtoType.Name}_Temp ({columns})";
}

// 生成仅包含主键的临时表(用于删除,优化性能)
private string GenerateCreateKeyOnlyTempTableSql(Type dtoType)
{
    // 用[Key]特性标记DTO主键,需提前在DTO类的主键属性上添加该特性
    var keyProp = dtoType.GetProperties()
                        .First(p => Attribute.IsDefined(p, typeof(KeyAttribute)));
    return $"CREATE TABLE #{dtoType.Name}_Temp ([{keyProp.Name}] {SqlTypeMappings[Nullable.GetUnderlyingType(keyProp.PropertyType) ?? keyProp.PropertyType]})";
}

3. DTO转DataTable工具方法

将DTO列表转换为DataTable,配合SqlBulkCopy导入临时表:

// 完整字段转换(用于新增、更新)
private DataTable ConvertToDataTable<T>(IEnumerable<T> items)
{
    var dt = new DataTable();
    var properties = typeof(T).GetProperties(BindingFlags.Public | BindingFlags.Instance);
    
    foreach (var prop in properties)
    {
        var columnType = Nullable.GetUnderlyingType(prop.PropertyType) ?? prop.PropertyType;
        dt.Columns.Add(prop.Name, columnType);
    }
    
    foreach (var item in items)
    {
        var row = dt.NewRow();
        foreach (var prop in properties)
        {
            row[prop.Name] = prop.GetValue(item) ?? DBNull.Value;
        }
        dt.Rows.Add(row);
    }
    return dt;
}

// 仅主键转换(用于删除)
private DataTable ConvertToKeyDataTable<T>(IEnumerable<T> items)
{
    var keyProp = typeof(T).GetProperties()
                        .First(p => Attribute.IsDefined(p, typeof(KeyAttribute)));
    var dt = new DataTable();
    dt.Columns.Add(keyProp.Name, Nullable.GetUnderlyingType(keyProp.PropertyType) ?? keyProp.PropertyType);
    
    foreach (var item in items)
    {
        dt.Rows.Add(keyProp.GetValue(item));
    }
    return dt;
}

4. 批量增删改SQL生成

批量插入SQL

private string GenerateBulkInsertSql(Type dtoType, string quotedSchema)
{
    var properties = dtoType.GetProperties(BindingFlags.Public | BindingFlags.Instance);
    var columnNames = string.Join(", ", properties.Select(p => $"[{p.Name}]"));
    var tempColumnNames = string.Join(", ", properties.Select(p => $"tmp.[{p.Name}]"));
    
    var keyProp = dtoType.GetProperties().First(p => Attribute.IsDefined(p, typeof(KeyAttribute)));
    
    return $@"
INSERT INTO {quotedSchema}.[{dtoType.Name}] ({columnNames})
SELECT {tempColumnNames}
FROM #{dtoType.Name}_Temp tmp
WHERE NOT EXISTS (
    SELECT 1 FROM {quotedSchema}.[{dtoType.Name}] t
    WHERE t.[{keyProp.Name}] = tmp.[{keyProp.Name}]
)";
}

批量更新SQL

private string GenerateBulkUpdateSql(Type dtoType, string quotedSchema)
{
    var keyProp = dtoType.GetProperties().First(p => Attribute.IsDefined(p, typeof(KeyAttribute)));
    var updateProps = dtoType.GetProperties(BindingFlags.Public | BindingFlags.Instance)
                            .Where(p => !Attribute.IsDefined(p, typeof(KeyAttribute)));
    var setClauses = string.Join(",\n", updateProps.Select(p => $"t.[{p.Name}] = tmp.[{p.Name}]"));
    
    return $@"
UPDATE t
SET {setClauses}
FROM {quotedSchema}.[{dtoType.Name}] t
JOIN #{dtoType.Name}_Temp tmp ON t.[{keyProp.Name}] = tmp.[{keyProp.Name}]";
}

批量删除SQL

private string GenerateBulkDeleteSql(Type dtoType, string quotedSchema)
{
    var keyProp = dtoType.GetProperties().First(p => Attribute.IsDefined(p, typeof(KeyAttribute)));
    
    return $@"
DELETE t
FROM {quotedSchema}.[{dtoType.Name}] t
JOIN #{dtoType.Name}_Temp tmp ON t.[{keyProp.Name}] = tmp.[{keyProp.Name}]";
}

5. 单事务内执行同步逻辑

public void Sync(string schemaName, Dictionary<Type, Dictionary<Synchronizable, TransactionType>> entitiesToSync)
{
    using var connection = new SqlConnection("你的数据库连接字符串");
    connection.Open();
    using var transaction = connection.BeginTransaction();

    try
    {
        // 安全处理Schema名称,避免SQL注入
        var quotedSchema = connection.ExecuteScalar<string>(
            "SELECT QUOTENAME(@schema)", 
            new { schema = schemaName }, 
            transaction: transaction
        );

        // 拆分新增、更新、删除分组(按DTO类型聚合)
        var addedGroups = entitiesToSync
            .SelectMany(kv => kv.Value.Where(e => e.Value == TransactionType.Added)
                                      .Select(e => (Type: kv.Key, Entity: e.Key)))
            .GroupBy(x => x.Type);
        
        var updatedGroups = entitiesToSync
            .SelectMany(kv => kv.Value.Where(e => e.Value == TransactionType.Updated)
                                      .Select(e => (Type: kv.Key, Entity: e.Key)))
            .GroupBy(x => x.Type);
        
        var deletedGroups = entitiesToSync
            .SelectMany(kv => kv.Value.Where(e => e.Value == TransactionType.Deleted)
                                      .Select(e => (Type: kv.Key, Entity: e.Key)))
            .GroupBy(x => x.Type);

        // 处理新增
        foreach (var group in addedGroups)
        {
            var dtoType = group.Key;
            var items = group.Select(x => x.Entity).Cast<object>();
            
            var createTempSql = GenerateCreateTempTableSql(dtoType);
            connection.Execute(createTempSql, transaction: transaction);
            
            var dt = ConvertToDataTable(items);
            using var bulkCopy = new SqlBulkCopy(connection, SqlBulkCopyOptions.Default, transaction);
            bulkCopy.DestinationTableName = $"#{dtoType.Name}_Temp";
            bulkCopy.WriteToServer(dt);
            
            var insertSql = GenerateBulkInsertSql(dtoType, quotedSchema);
            connection.Execute(insertSql, transaction: transaction);
            
            connection.Execute($"DROP TABLE #{dtoType.Name}_Temp", transaction: transaction);
        }

        // 处理更新
        foreach (var group in updatedGroups)
        {
            var dtoType = group.Key;
            var items = group.Select(x => x.Entity).Cast<object>();
            
            var createTempSql = GenerateCreateTempTableSql(dtoType);
            connection.Execute(createTempSql, transaction: transaction);
            
            var dt = ConvertToDataTable(items);
            using var bulkCopy = new SqlBulkCopy(connection, SqlBulkCopyOptions.Default, transaction);
            bulkCopy.DestinationTableName = $"#{dtoType.Name}_Temp";
            bulkCopy.WriteToServer(dt);
            
            var updateSql = GenerateBulkUpdateSql(dtoType, quotedSchema);
            connection.Execute(updateSql, transaction: transaction);
            
            connection.Execute($"DROP TABLE #{dtoType.Name}_Temp", transaction: transaction);
        }

        // 处理删除
        foreach (var group in deletedGroups)
        {
            var dtoType = group.Key;
            var items = group.Select(x => x.Entity).Cast<object>();
            
            var createTempSql = GenerateCreateKeyOnlyTempTableSql(dtoType);
            connection.Execute(createTempSql, transaction: transaction);
            
            var dt = ConvertToKeyDataTable(items);
            using var bulkCopy = new SqlBulkCopy(connection, SqlBulkCopyOptions.Default, transaction);
            bulkCopy.DestinationTableName = $"#{dtoType.Name}_Temp";
            bulkCopy.WriteToServer(dt);
            
            var deleteSql = GenerateBulkDeleteSql(dtoType, quotedSchema);
            connection.Execute(deleteSql, transaction: transaction);
            
            connection.Execute($"DROP TABLE #{dtoType.Name}_Temp", transaction: transaction);
        }

        transaction.Commit();
    }
    catch
    {
        transaction.Rollback();
        throw;
    }
}

关键注意事项

  • 防注入保障:所有表名、列名都用方括号包裹,Schema名称通过QUOTENAME函数处理,彻底避免拼接用户输入导致的注入风险。
  • 主键识别:通过[Key]特性标记DTO的主键属性,替代硬编码字段,提升方案通用性。
  • 性能优化:SqlBulkCopy是SQL Server批量导入的最优方案,比生成多值INSERT语句效率高数倍,适合数千条数据的场景。
  • 事务一致性:整个同步流程在单个事务内执行,任意步骤失败都会回滚所有操作,保证数据完整性。

内容的提问来源于stack exchange,提问作者Zenith

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 06:25:55