/
rezvich
/
SoapBatchProcessor
Обзор
Документация
Войти
/
rezvich
/
SoapBatchProcessor
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
Data/BulkWriter.cs
286 строк
11 KB
rezvich
Подправил токен
04 авг 2026, 12:20
04 авг 2026, 12:20
632636d
Код
Авторство
О чём код?
using Microsoft.Data.SqlClient; using SoapBatchProcessor.Models.Data_Models; using System.Data; namespace SoapBatchProcessor.Data { public class BulkWriter { private readonly string _cs; private readonly int _batchSize; private readonly int _timeout; public BulkWriter(string cs, int batchSize, int timeout) { _cs = cs; _batchSize = batchSize; _timeout = timeout; } public async Task<HashSet<long>> FlushAsync(Accumulator acc) { var failed = new HashSet<long>(); using var con = new SqlConnection(_cs); await con.OpenAsync(); using (var cmd = new SqlCommand("SET XACT_ABORT ON;", con)) await cmd.ExecuteNonQueryAsync(); // PersonData (0..N на ответ) { var rows = acc.Responses .Where(r => r.Person != null) .Select(r => (r.Person!, r.ResponseId, r.RequestId)) .ToList(); if (rows.Count > 0) { var dt = ToDataTableWithParent("PersonData", rows); var f = await SafeBulkAsync(con, dt); foreach (var id in f) failed.Add(id); } } // OmsPolicyData { var rows = acc.Responses .SelectMany(r => (r.OmsPolicies ?? new List<OmsPolicyData>()) .Select(x => (item: x, rid: r.ResponseId, reqId: r.RequestId))) .ToList(); if (rows.Count > 0) { var dt = ToDataTableWithParent("OmsPolicyData", rows); var f = await SafeBulkAsync(con, dt); foreach (var id in f) failed.Add(id); } } // DudlData { var rows = acc.Responses .SelectMany(r => (r.Dudls ?? new List<DudlData>()) .Select(x => (x, r.ResponseId, r.RequestId))) .ToList(); if (rows.Count > 0) { var dt = ToDataTableWithParent("DudlData", rows); var f = await SafeBulkAsync(con, dt); foreach (var id in f) failed.Add(id); } } // AddressData { var rows = acc.Responses .SelectMany(r => (r.Addresses ?? new List<AddressData>()) .Select(x => (item: x, rid: r.ResponseId, reqId: r.RequestId))) .ToList(); if (rows.Count > 0) { var dt = ToDataTableWithParent("AddressData", rows); var f = await SafeBulkAsync(con, dt); foreach (var id in f) failed.Add(id); } } // AttachData { var rows = acc.Responses .SelectMany(r => (r.Attachments ?? new List<AttachData>()) .Select(x => (item: x, rid: r.ResponseId, reqId: r.RequestId))) .ToList(); if (rows.Count > 0) { var dt = ToDataTableWithParent("AttachData", rows); var f = await SafeBulkAsync(con, dt); foreach (var id in f) failed.Add(id); } } // ContactData { var rows = acc.Responses .SelectMany(r => (r.Contacts ?? new List<ContactData>()) .Select(x => (item: x, rid: r.ResponseId, reqId: r.RequestId))) .ToList(); if (rows.Count > 0) { var dt = ToDataTableWithParent("ContactData", rows); var f = await SafeBulkAsync(con, dt); foreach (var id in f) failed.Add(id); } } // SnilsData (по одному на ответ, но держим общий паттерн) { var rows = acc.Responses .Where(r => r.Snils != null) .Select(r => (item: r.Snils!, rid: r.ResponseId, reqId: r.RequestId)) .ToList(); if (rows.Count > 0) { var dt = ToDataTableWithParent("SnilsData", rows); var f = await SafeBulkAsync(con, dt); foreach (var id in f) failed.Add(id); } } // SocialStatusData { var rows = acc.Responses .SelectMany(r => (r.SocialStatuses ?? new List<SocialStatusData>()) .Select(x => (item: x, rid: r.ResponseId, reqId: r.RequestId))) .ToList(); if (rows.Count > 0) { var dt = ToDataTableWithParent("SocialStatusData", rows); var f = await SafeBulkAsync(con, dt); foreach (var id in f) failed.Add(id); } } // ErnData { var rows = acc.Responses .Where(r => r.Ern != null) .Select(r => (item: r.Ern!, rid: r.ResponseId, reqId: r.RequestId)) .ToList(); if (rows.Count > 0) { var dt = ToDataTableWithParent("ErnData", rows); var f = await SafeBulkAsync(con, dt); foreach (var id in f) failed.Add(id); } } return failed; } private async Task<HashSet<long>> SafeBulkAsync(SqlConnection con, DataTable dt) { var failed = new HashSet<long>(); // 1) Пытаемся целиком try { using var bc = new SqlBulkCopy(con, SqlBulkCopyOptions.CheckConstraints | SqlBulkCopyOptions.UseInternalTransaction, null) { DestinationTableName = $"dbo.{dt.TableName}", BatchSize = _batchSize, BulkCopyTimeout = _timeout }; bc.ColumnMappings.Clear(); foreach (DataColumn c in dt.Columns) bc.ColumnMappings.Add(c.ColumnName, c.ColumnName); await bc.WriteToServerAsync(dt); Console.WriteLine($"[BULK OK] {dt.TableName}: {dt.Rows.Count} rows"); return failed; } catch (Exception ex) { Console.Error.WriteLine($"[BULK WARN] Table={dt.TableName}: {ex.Message} → fallback to per-row INSERT"); } // 2) Фолбэк: параметризованный INSERT по одной строке // — соберём один SqlCommand с параметрами, будем переиспользовать var cols = dt.Columns.Cast<DataColumn>().ToList(); var colNames = string.Join(",", cols.Select(c => $"[{c.ColumnName}]")); var paramNames = string.Join(",", cols.Select(c => $"@{c.ColumnName}")); var sql = $"INSERT INTO dbo.{dt.TableName} ({colNames}) VALUES ({paramNames});"; using var cmd = new SqlCommand(sql, con) { CommandTimeout = _timeout }; cmd.Parameters.Clear(); foreach (var c in cols) cmd.Parameters.Add(new SqlParameter($"@{c.ColumnName}", DBNull.Value)); foreach (DataRow r in dt.Rows) { try { // мостик типов для известных колонок if (dt.TableName == "DudlData" && dt.Columns.Contains("NoCitizenship")) { var v = r["NoCitizenship"]; if (v is bool b) r["NoCitizenship"] = b; // если в БД bit // если всё ещё nvarchar в БД, то: // r["NoCitizenship"] = b ? "true" : "false"; } foreach (var c in cols) { var val = r[c]; cmd.Parameters[$"@{c.ColumnName}"].Value = (val == null || val is DBNull) ? DBNull.Value : val; } await cmd.ExecuteNonQueryAsync(); } catch (Exception exRow) { long reqId = 0; try { reqId = Convert.ToInt64(r["RequestId"]); } catch { /* ignore */ } failed.Add(reqId); Console.Error.WriteLine($"[ROW FAIL] Table={dt.TableName}, RequestId={reqId}: {exRow.Message}"); } } Console.WriteLine($"[BULK PARTIAL] {dt.TableName}: ok={dt.Rows.Count - failed.Count}, failed={failed.Count}"); return failed; } private static DataTable ToDataTableWithParent<T>(string tableName, IEnumerable<(T item, long responseId, long requestId)> items) { var dt = new DataTable(tableName); dt.Columns.Add("ResponseId", typeof(long)); dt.Columns.Add("RequestId", typeof(long)); var props = typeof(T).GetProperties(); foreach (var p in props) { var t = Nullable.GetUnderlyingType(p.PropertyType) ?? p.PropertyType; dt.Columns.Add(p.Name, t); } foreach (var (item, rid, reqId) in items) { var row = dt.NewRow(); row["ResponseId"] = rid; row["RequestId"] = reqId; foreach (var p in props) row[p.Name] = p.GetValue(item) ?? DBNull.Value; dt.Rows.Add(row); } return dt; } private async Task Bulk(DataTable dt, SqlConnection con) { using var bc = new SqlBulkCopy(con) { DestinationTableName = $"dbo.{dt.TableName}", BulkCopyTimeout = _timeout, BatchSize = _batchSize }; foreach (DataColumn c in dt.Columns) bc.ColumnMappings.Add(c.ColumnName, c.ColumnName); await bc.WriteToServerAsync(dt); bc.NotifyAfter = _batchSize; // будет дёргать событие раз в батч bc.SqlRowsCopied += (_, e) => { Console.WriteLine($"[BULK] {dt.TableName}: {e.RowsCopied} rows copied..."); }; } } }