基于C# Dapper的单事务批量增删改动态SQL实现方案
基于Dapper的单事务批量同步解决方案
针对你的需求,核心解决方案是临时表+SqlBulkCopy+批量SQL操作,既满足单事务要求,又能规避参数数量限制,同时通过动态生成安全的SQL实现防注入,完全不需要依赖EF。
核心思路
- 利用Session级临时表批量导入待同步数据,替代逐个生成参数的方式,避开SQL Server参数数量上限(默认2100个)。
- 通过反射动态生成临时表结构、目标表的增删改SQL,保证DTO与数据库表的一一对应。
- 所有操作在单个数据库事务内执行,确保同步的原子性。
- 采用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
相关产品推荐
相关产品推荐

