【问题标题】:Export a large CSV file in parallel to SQL server将大型 CSV 文件并行导出到 SQL Server
【发布时间】:2014-12-17 11:51:15
【问题描述】:

我有一个大的 CSV 文件……我的硬盘上有 10 列、1 亿行、大约 6 GB 大小。 我想逐行读取这个 CSV 文件,然后使用 SQL 大容量复制将数据加载到 Microsoft SQL 服务器数据库中。 我在这里和互联网上阅读了几个主题。大多数人认为,并行读取 CSV 文件并不能提高效率,因为任务/线程会争用磁盘访问权限。

我想要做的是,从 CSV 中逐行读取并将其添加到阻止大小为 100K 行的集合中。一旦这个集合完全启动一个新的任务/线程,使用 SQLBuckCopy API 将数据写入 SQL 服务器。

我已经编写了这段代码,但在运行时遇到了一个错误,提示“尝试在具有挂起操作的对象上调用大容量复制”。这种情况看起来可以使用 .NET 4.0 TPL 轻松解决,但我无法让它工作。关于我做错了什么有什么建议吗?

    public static void LoadCsvDataInParalleToSqlServer(string fileName, string connectionString, string table, DataColumn[] columns, bool truncate)
    {
        const int inputCollectionBufferSize = 1000000;
        const int bulkInsertBufferCapacity = 100000;
        const int bulkInsertConcurrency = 8;

        var sqlConnection = new SqlConnection(connectionString);
        sqlConnection.Open();

        var sqlBulkCopy = new SqlBulkCopy(sqlConnection.ConnectionString, SqlBulkCopyOptions.TableLock)
        {
            EnableStreaming = true,
            BatchSize = bulkInsertBufferCapacity,
            DestinationTableName = table,
            BulkCopyTimeout = (24 * 60 * 60),
        };

        BlockingCollection<DataRow> rows = new BlockingCollection<DataRow>(inputCollectionBufferSize);
        DataTable dataTable = new DataTable(table);
        dataTable.Columns.AddRange(columns);

        Task loadTask = Task.Factory.StartNew(() =>
            {
                foreach (DataRow row in ReadRows(fileName, dataTable))
                {
                    rows.Add(row);
                }

                rows.CompleteAdding();
            });

        List<Task> insertTasks = new List<Task>(bulkInsertConcurrency);

        for (int i = 0; i < bulkInsertConcurrency; i++)
        {
            insertTasks.Add(Task.Factory.StartNew((x) =>
                {
                    List<DataRow> bulkInsertBuffer = new List<DataRow>(bulkInsertBufferCapacity);

                    foreach (DataRow row in rows.GetConsumingEnumerable())
                    {
                        if (bulkInsertBuffer.Count == bulkInsertBufferCapacity)
                        {
                            SqlBulkCopy bulkCopy = x as SqlBulkCopy;
                            var dataRows = bulkInsertBuffer.ToArray();
                            bulkCopy.WriteToServer(dataRows);
                            Console.WriteLine("Inserted rows " + bulkInsertBuffer.Count);
                            bulkInsertBuffer.Clear();
                        }

                        bulkInsertBuffer.Add(row);
                    }

                },
                sqlBulkCopy));
        }

        loadTask.Wait();
        Task.WaitAll(insertTasks.ToArray());
    }

    private static IEnumerable<DataRow> ReadRows(string fileName, DataTable dataTable)
    {
        using (var textFieldParser = new TextFieldParser(fileName))
        {
            textFieldParser.TextFieldType = FieldType.Delimited;
            textFieldParser.Delimiters = new[] { "," };
            textFieldParser.HasFieldsEnclosedInQuotes = true;

            while (!textFieldParser.EndOfData)
            {
                string[] cols = textFieldParser.ReadFields();

                DataRow row = dataTable.NewRow();

                for (int i = 0; i < cols.Length; i++)
                {
                    if (string.IsNullOrEmpty(cols[i]))
                    {
                        row[i] = DBNull.Value;
                    }
                    else
                    {
                        row[i] = cols[i];
                    }
                }

                yield return row;
            }
        }
    }

【问题讨论】:

  • 与其花时间编写自己的工具,不如使用已经完成此任务的 ETL 工具,例如 SQL Server Integration Services。
  • 您是否尝试过此代码的顺序版本并证明多线程的复杂性值得性能提升?
  • 有很多优化批量插入的在线指南,即technet.microsoft.com/en-us/library/ms190421(v=sql.105).aspx。听起来您正在尝试解决尚未证明存在的问题。我建议你先简单地使用BCP.EXE 获得一个基线,然后尝试改进那个时间。
  • 根据我在网上阅读的内容...SqlBulkCopy 比 SQL Server 拥有的内置数据导入工具快得多,我相信它在幕后使用了 SSIS。负载性能至关重要,因此我对为其编写自己的 lil 应用程序进行了调查
  • 我有类似的卷,在我的情况下,我的 SQL 服务器的磁盘 IO 是瓶颈,所以我确实拆分了批次,但我没有并行。

标签: c# sql sql-server multithreading csv


【解决方案1】:

不要。

并行访问可能会或可能不会让您更快地读取文件(它不会,但我不会打战斗...)但对于某些并行写入它会赢'不要给你更快的批量插入。这是因为最少记录的批量插入(即非常快批量插入)需要一个表锁。见Prerequisites for Minimal Logging in Bulk Import

最小日志记录要求目标表满足以下条件:

...
- 已指定表锁定(使用 TABLOCK)
...

根据定义,并行插入无法获得并发表锁。 QED。你找错树了。

停止从互联网上随机找到您的资源。阅读The Data Loading Performance Guide 是...高性能数据加载指南。

我会建议你停止发明轮子。使用SSIS,这正是旨在处理的内容。

【讨论】:

  • 好的。你能指出一个现有的 SSIS 包,它可以将 csv 文件中的数据批量插入到 SQL 表中吗?我不想像你说的那样重新发明轮子,所以想使用一些预先存在的解决方案,而不是自己创建 ssis 包
  • 您只需要一个Flat File Source 连接到具有快速加载设置的OleDB Destination。例如,请参阅Import CSV File into Database Table Using SSIS
  • >> 但对于某些并行写入,它不会为您提供更快的批量插入。
【解决方案2】:

http://joshclose.github.io/CsvHelper/

https://efbulkinsert.codeplex.com/

如果可能的话,我建议您使用前面提到的 csvhelper 将文件读入 List 并像您正在做的那样使用批量插入或我使用过且速度惊人的 efbulkinsert 写入您的数据库。

using CsvHelper;

public static List<T> CSVImport<T,TClassMap>(string csvData, bool hasHeaderRow, char delimiter, out string errorMsg) where TClassMap : CsvHelper.Configuration.CsvClassMap
    {
        errorMsg = string.Empty;
        var result = Enumerable.Empty<T>();

        MemoryStream memStream = new MemoryStream(Encoding.UTF8.GetBytes(csvData));
        StreamReader streamReader = new StreamReader(memStream);
        var csvReader = new CsvReader(streamReader);

        csvReader.Configuration.RegisterClassMap<TClassMap>();
        csvReader.Configuration.DetectColumnCountChanges = true;
        csvReader.Configuration.IsHeaderCaseSensitive = false;
        csvReader.Configuration.TrimHeaders = true;
        csvReader.Configuration.Delimiter = delimiter.ToString();
        csvReader.Configuration.SkipEmptyRecords = true;
        List<T> items = new List<T>();

        try
        {
            items = csvReader.GetRecords<T>().ToList();
        }
        catch (Exception ex)
        {
            while (ex != null)
            {
                errorMsg += ex.Message + Environment.NewLine;

                foreach (var val in ex.Data.Values)
                    errorMsg += val.ToString() + Environment.NewLine;

                ex = ex.InnerException;
            }
        }
        return items;
    }
}

编辑 - 我不明白您在使用批量插入做什么。您要批量插入整个列表或数据数据表,而不是逐行插入。

【讨论】:

  • csv 似乎很大。 (6 GB)。 GetRecords&lt;T&gt;().ToList() 是否将所有内容加载到内存中?
  • 是的 - 好点,这对他来说可能是不可能的。 BulkInserting in one gulp 是一个很大的节省时间。或许他可以调用列表中的 Take() 来将其分块。他的列表似乎适合 File.ReadLines 创建的字符串。
  • 由于大小,我无法将整个 csv 文件加载到内存中。所以我需要一次读取 100K 行,然后使用 bulkinsert 将其写入 SQL Server。是的,我确实想一次写整个表格,而不是一次写一行。
【解决方案3】:

您可以创建存储过程并传递文件位置,如下所示

CREATE PROCEDURE [dbo].[CSVReaderTransaction]
    @Filepath varchar(100)=''
AS
-- STEP 1: Start the transaction
BEGIN TRANSACTION

-- STEP 2 & 3: checking @@ERROR after each statement
EXEC ('BULK INSERT Employee FROM ''' +@Filepath
        +''' WITH (FIELDTERMINATOR = '','', ROWTERMINATOR = ''\n'' )')

-- Rollback the transaction if there were any errors
IF @@ERROR <> 0
 BEGIN
    -- Rollback the transaction
    ROLLBACK

    -- Raise an error and return
    RAISERROR ('Error in inserting data into employee Table.', 16, 1)
    RETURN
 END

COMMIT TRANSACTION

您还可以添加 BATCHSIZE 选项,例如 FIELDTERMINATOR 和 ROWTERMINATOR。

【讨论】:

  • 对我不起作用,因为该文件位于我的本地计算机上,而 SQL 服务器位于另一台计算机上,并且该计算机无法远程访问我的本地驱动器
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2013-02-20
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-06-08
  • 1970-01-01
相关资源
最近更新 更多