diff --git a/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs b/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs index 8713b7e..16d4a4f 100644 --- a/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs +++ b/src/Predictalytics.Infrastructure/Data/Repositories/TradeRepository.cs @@ -48,7 +48,10 @@ public class TradeRepository : ITradeRepository public async Task AddRangeAsync(IEnumerable trades, CancellationToken ct = default) { - foreach (var t in trades) + var tradeList = trades.ToList(); + if (tradeList.Count == 0) return; + + foreach (var t in tradeList) { t.Outcome = StringHelper.Truncate(t.Outcome, 128) ?? ""; t.PlatformTradeId = StringHelper.Truncate(t.PlatformTradeId, 256) ?? ""; @@ -58,31 +61,38 @@ public class TradeRepository : ITradeRepository t.TransactionHash = StringHelper.Truncate(t.TransactionHash, 66); } - try + foreach (var chunk in tradeList.Chunk(1000)) { - _db.Trades.AddRange(trades); - await _db.SaveChangesAsync(ct); - } - catch - { - foreach (var t in trades) - { - try { _db.Entry(t).State = EntityState.Detached; } catch { } - } + var sb = new System.Text.StringBuilder("INSERT IGNORE INTO Trades (PlatformTradeId, MarketId, AssetId, Outcome, Side, Price, Size, Amount, ExecutedAt, TransactionHash, TraderId, MarketOutcomeId, DbMarketId, Platform, IsContextEnriched) VALUES "); + var parameters = new List(); - // Fallback to inserting one by one to avoid losing the entire batch due to a single duplicate - foreach (var t in trades) + for (int i = 0; i < chunk.Length; i++) { - try - { - _db.Trades.Add(t); - await _db.SaveChangesAsync(ct); - } - catch - { - try { _db.Entry(t).State = EntityState.Detached; } catch { } - } + var t = chunk[i]; + int pIdx = i * 15; + sb.Append($"({{{pIdx}}}, {{{pIdx + 1}}}, {{{pIdx + 2}}}, {{{pIdx + 3}}}, {{{pIdx + 4}}}, {{{pIdx + 5}}}, {{{pIdx + 6}}}, {{{pIdx + 7}}}, {{{pIdx + 8}}}, {{{pIdx + 9}}}, {{{pIdx + 10}}}, {{{pIdx + 11}}}, {{{pIdx + 12}}}, {{{pIdx + 13}}}, {{{pIdx + 14}}})"); + + if (i < chunk.Length - 1) + sb.Append(", "); + + parameters.Add(t.PlatformTradeId); + parameters.Add(t.MarketId); + parameters.Add(t.AssetId); + parameters.Add(t.Outcome); + parameters.Add((int)t.Side); + parameters.Add(t.Price); + parameters.Add(t.Size); + parameters.Add(t.Amount); + parameters.Add(t.ExecutedAt); + parameters.Add(t.TransactionHash ?? (object)DBNull.Value); + parameters.Add(t.TraderId); + parameters.Add(t.MarketOutcomeId ?? (object)DBNull.Value); + parameters.Add(t.DbMarketId ?? (object)DBNull.Value); + parameters.Add((int)t.Platform); + parameters.Add(t.IsContextEnriched); } + + await _db.Database.ExecuteSqlRawAsync(sb.ToString(), parameters.ToArray(), ct); } }