using MoreLinq; using MySqlConnector; using Sony.Filtr.Database; using Sony.Filtr.Functional; using System; using System.Collections.Generic; using System.IO; using System.Linq; using System.Text; using System.Threading.Tasks; namespace Sony.Filtr.Playlists { public class BulkLoader { private const int MaxRowsPerFile = 20_000; public static async Task LoadAsync(string tableName, Func[] fields, string[] columns, IEnumerable items, MySqlBulkLoaderConflictOption conflictOptions = MySqlBulkLoaderConflictOption.Ignore) { using (var connection = await DatabaseHandler.GetOpenConnectionAsync()) { var transaction = connection.BeginTransaction(); foreach (var filePath in GenerateFileBatches(items, fields)) { try { await LoadBatch(filePath, connection, tableName, columns, conflictOptions); } catch (Exception ex) { transaction.Rollback(); throw; } finally { if (File.Exists(filePath)) { File.Delete(filePath); } } } transaction.Commit(); } } public static Func> ToExtensionFunc(string tableName, Func[] fields, string[] columns, IEnumerable items, MySqlBulkLoaderConflictOption conflictOptions = MySqlBulkLoaderConflictOption.Ignore) { return FuncUtils.Partial[], string[], IEnumerable, MySqlBulkLoaderConflictOption, Task>( LoadAsync, tableName) .Partial(fields) .Partial(columns) .Partial(items) .ToUnit(); } private static async Task LoadBatch(string filePath, MySqlConnection connection, string tableName, string[] columns, MySqlBulkLoaderConflictOption conflictOptions) { var bulkLoader = new MySqlBulkLoader(connection); bulkLoader.Local = true; bulkLoader.FileName = filePath; bulkLoader.TableName = tableName; bulkLoader.FieldTerminator = "|"; bulkLoader.LineTerminator = Environment.NewLine; bulkLoader.ConflictOption = conflictOptions; bulkLoader.CharacterSet = "utf8mb4"; bulkLoader.Columns.AddRange(columns); await bulkLoader.LoadAsync(); } private static string GenerateFileName() { return $"{typeof(T).Name}_{Guid.NewGuid().ToString()}.csv"; } //private IEnumerable GenerateFileBatchesAsync() //{ // foreach (var batch in this.Items.Batch(MaxRowsPerFile)) // { // string filePath = GetFileName(); // using (var saveFile = File.Open(filePath, FileMode.OpenOrCreate)) // { // using (var fileWriter = new StreamWriter(saveFile) { AutoFlush = false }) // { // foreach (var item in batch) // { // //await fileWriter.WriteLineAsync(ItemToString(item, "|")); // fileWriter.WriteLine(this.ItemToString(item, "|")); // } // } // } // yield return filePath; // } //} private static IEnumerable GenerateFileBatches(IEnumerable items, Func[] fields) { foreach (var batch in items.Batch(MaxRowsPerFile)) { string filePath = GenerateFileName(); using (var saveFile = File.Open(filePath, FileMode.OpenOrCreate)) { using (var fileWriter = new StreamWriter(saveFile) { AutoFlush = false }) { StringBuilder builder = new StringBuilder(); foreach (var item in batch) { builder.AppendLine(ItemToString(item, "|", fields)); } fileWriter.Write(builder.ToString()); } } yield return filePath; } } private static string ItemToString(T item, string separator, Func[] fields) { return String.Join(separator, fields.Select(fieldGetter => fieldGetter(item))); } } }