C#:用Sylvan.Data.CSV实现SQL的CSV行新增与更新(解决BulkCopy报错)
CSV数据批量同步到SQL:实现Upsert(更新/插入)
问题根源
你碰到的报错是因为SqlBulkCopy.WriteToServer根本不支持单条行的写入逻辑——你用csv.Read()定位到当前行后,直接把reader传给WriteToServer,它会尝试从当前位置读到CSV末尾,而且逐行调用批量复制完全浪费了它的性能优势,写法本身就不符合API的设计逻辑。
解决方案一:逐行处理(小数据量适用)
如果你的CSV数据量很小(比如几百条),可以调整插入逻辑,改用参数化语句处理单行:
static void LoadTableCsv(SqlConnection conn, string tableName, string csvFile) { // 读取目标表的列结构 var cmd = conn.CreateCommand(); conn.Open(); cmd.CommandText = $"select top 0 * from {tableName}"; // 注意:若表名来自用户输入,必须做白名单校验防注入 var reader = cmd.ExecuteReader(); var colSchema = reader.GetColumnSchema(); reader.Close(); // 把SQL表结构应用到CSV读取器,自动处理类型映射 var csvSchema = new CsvSchema(colSchema); var csvOpts = new CsvDataReaderOptions { Schema = csvSchema }; using var csv = CsvDataReader.Create(csvFile, csvOpts); // 初始化记录存在性检查命令 using var checkCommand = new SqlCommand($"SELECT COUNT(*) FROM {tableName} WHERE TicketID = @value", conn); checkCommand.Parameters.Add("@value", SqlDbType.Int); // 预编译参数化插入命令 var insertColList = string.Join(", ", colSchema.Select(c => c.ColumnName)); var insertParamList = string.Join(", ", colSchema.Select(c => $"@{c.ColumnName}")); using var insertCommand = new SqlCommand($"INSERT INTO {tableName} ({insertColList}) VALUES ({insertParamList})", conn); // 批量添加插入参数 foreach (var col in colSchema) { var sqlType = GetSqlDbType(col.DataType); insertCommand.Parameters.Add($"@{col.ColumnName}", sqlType); } // 预编译参数化更新命令 var updateSetClause = string.Join(", ", colSchema.Where(c => c.ColumnName != "TicketID") .Select(c => $"{c.ColumnName} = @{c.ColumnName}")); using var updateCommand = new SqlCommand($"UPDATE {tableName} SET {updateSetClause} WHERE TicketID = @TicketID", conn); // 批量添加更新参数 foreach (var col in colSchema) { var sqlType = GetSqlDbType(col.DataType); updateCommand.Parameters.Add($"@{col.ColumnName}", sqlType); } // 遍历CSV每一行 while (csv.Read()) { var ticketId = csv.GetInt32(csv.GetOrdinal("TicketID")); checkCommand.Parameters["@value"].Value = ticketId; bool recordExists = (int)checkCommand.ExecuteScalar() > 0; if (recordExists) { // 更新现有记录:为每个参数赋值 foreach (var col in colSchema) { int ordinal = csv.GetOrdinal(col.ColumnName); updateCommand.Parameters[$"@{col.ColumnName}"].Value = csv.IsDBNull(ordinal) ? DBNull.Value : csv.GetValue(ordinal); } updateCommand.ExecuteNonQuery(); } else { // 插入新记录:为每个参数赋值 foreach (var col in colSchema) { int ordinal = csv.GetOrdinal(col.ColumnName); insertCommand.Parameters[$"@{col.ColumnName}"].Value = csv.IsDBNull(ordinal) ? DBNull.Value : csv.GetValue(ordinal); } insertCommand.ExecuteNonQuery(); } } conn.Close(); } // 辅助方法:将CLR类型转换为对应SqlDbType private static SqlDbType GetSqlDbType(Type clrType) { if (clrType == typeof(int)) return SqlDbType.Int; if (clrType == typeof(string)) return SqlDbType.NVarChar; if (clrType == typeof(DateTime)) return SqlDbType.DateTime; // 根据你的表结构补充其他类型映射 throw new NotSupportedException($"不支持的类型:{clrType.Name}"); }
解决方案二:批量Upsert(大数据量必用)
逐行操作在数据量大时性能极差,推荐用临时表+MERGE语句的方案,这是SQL Server批量同步数据的最优解:
static void LoadTableCsvBulk(SqlConnection conn, string tableName, string csvFile) { conn.Open(); // 1. 创建与目标表结构完全一致的临时表 string tempTableName = $"#{tableName}_Temp"; using var createTempTableCmd = new SqlCommand($"SELECT TOP 0 * INTO {tempTableName} FROM {tableName}", conn); createTempTableCmd.ExecuteNonQuery(); // 2. 读取目标表结构并应用到CSV读取器 var cmd = conn.CreateCommand(); cmd.CommandText = $"select top 0 * from {tableName}"; var reader = cmd.ExecuteReader(); var colSchema = reader.GetColumnSchema(); reader.Close(); var csvSchema = new CsvSchema(colSchema); var csvOpts = new CsvDataReaderOptions { Schema = csvSchema }; using var csv = CsvDataReader.Create(csvFile, csvOpts); // 3. 用SqlBulkCopy批量导入CSV数据到临时表 using var bulkCopy = new SqlBulkCopy(conn); bulkCopy.DestinationTableName = tempTableName; bulkCopy.EnableStreaming = true; bulkCopy.WriteToServer(csv); // 4. 用MERGE语句完成批量Upsert:存在则更新,不存在则插入 string mergeColList = string.Join(", ", colSchema.Select(c => c.ColumnName)); string updateSetClause = string.Join(", ", colSchema.Where(c => c.ColumnName != "TicketID") .Select(c => $"Target.{c.ColumnName} = Source.{c.ColumnName}")); string mergeSql = $@" MERGE INTO {tableName} AS Target USING {tempTableName} AS Source ON Target.TicketID = Source.TicketID WHEN MATCHED THEN UPDATE SET {updateSetClause} WHEN NOT MATCHED THEN INSERT ({mergeColList}) VALUES ({string.Join(", ", colSchema.Select(c => $"Source.{c.ColumnName}"))}); -- 清理临时表 DROP TABLE {tempTableName}; "; using var mergeCmd = new SqlCommand(mergeSql, conn); mergeCmd.ExecuteNonQuery(); conn.Close(); }
重要提示
- 防SQL注入:如果
tableName来自用户输入,必须做白名单校验,绝对不能直接拼接SQL; - 性能差异:方案二的批量操作性能比逐行处理高几十甚至上百倍,数据量越大优势越明显;
- 类型安全:借助Sylvan.Data.CSV的Schema映射,自动处理CSV到SQL的类型转换,避免手动解析字符串的错误;
- 资源管理:所有实现
IDisposable的对象(如SqlCommand、SqlDataReader)都用using包裹,确保资源正确释放。
内容的提问来源于stack exchange,提问作者NightBladium
相关产品推荐
相关产品推荐

