【问题标题】:Inserting multiple rows from ConcurrentQueue into SQL Server table将 ConcurrentQueue 中的多行插入 SQL Server 表
【发布时间】:2014-10-27 14:49:17
【问题描述】:

我有一个 c# 应用程序,它从一个数据库中获取数据,进行必要的转换,并将数据插入到另一个数据库的表中。为此,我将源数据插入队列,然后处理队列以将数据插入目标表。我有两个单独的线程来读取源数据和写入目标数据。读取线程的运行速度比写入线程快得多,所以我的队列很快就会被填满。

正如您在阅读线程中看到的,我正在使用 SqlCommand.ExecuteReader() 来读取数据。然后我遍历队列并为每一行执行单独的 INSERT 语句。我所设想的是,不是逐行插入,而是做某种(可能是 Linq 语句)。有谁知道如何更快地完成我的插入操作?

队列定义:

private static readonly BlockingCollection<HistorianData> ValueQueue = new BlockingCollection<HistorianData>(new ConcurrentQueue<HistorianData>(), 1000000);

阅读:

    public static void EnqueueHistorianData(SqlConnection connection, int idToAdd, DateTime minDatetime, DateTime maxDateTime, string cluster, string dbName, string dataTable, string idTable, string mainIdColumn, string foreignIdColumn, string dateColumn, string nameColumn, string valueColumn)
    {
        StringBuilder select = new StringBuilder();
        HistorianData values;

        select.Append(String.Format("SELECT {0}.{1},", dataTable, dateColumn));
        select.Append(String.Format(" '{0}.' + {1}.{2},", cluster, idTable, nameColumn));
        select.Append(String.Format(" {0}.{1}", dataTable, valueColumn));
        select.Append(String.Format(" FROM {0}.{1}", dbName, dataTable));
        select.Append(String.Format(" JOIN {0}.{1}", dbName, idTable));
        select.Append(String.Format(" ON {0}.{1} = {2}.{3}", dataTable, foreignIdColumn, idTable, mainIdColumn));
        select.Append(String.Format(" INNER JOIN Runtime.dbo.Tag"));
        select.Append(String.Format(" ON Runtime.dbo.Tag.TagName = '{0}.' + {1}.{2}", cluster, idTable, nameColumn));
        select.Append(String.Format(" WHERE {0}.{1} >= '{2}'", dataTable, dateColumn, minDatetime.ToString()));
        select.Append(String.Format(" AND {0}.{1} = {2}", dataTable, foreignIdColumn, idToAdd.ToString()));
        select.Append(String.Format(" AND {0}.{1} >= '{2}'", dataTable, dateColumn, minDatetime.ToString()));
        select.Append(String.Format(" AND {0}.{1} < '{2}'", dataTable, dateColumn, maxDateTime.ToString()));

        using (var command = new SqlCommand(select.ToString(), connection))
        {
            command.CommandTimeout = 1000;

            using (var reader = command.ExecuteReader())
            {
                if (reader.HasRows)
                {
                    while (reader.Read())
                    {
                        values = new HistorianData();

                        values.SampleDate = reader.GetDateTime(0);
                        values.TagName = reader.GetString(1);
                        values.TagValue = reader.GetDouble(2);

                        ValueQueue.Add(values);
                    }

                    values = null;
                    reader.Close();
                }
            }
        }
    }

写作:

    public static void WriteQueueValuesToHistorian(string connectionString)
    {   
        HistorianData values;

        using (var connection = new SqlConnection(connectionString))
        {
            connection.Open();

            using (SqlCommand insertCommand = connection.CreateCommand())
            {
                insertCommand.CommandType = CommandType.Text;
                insertCommand.CommandText = "INSERT INTO Runtime.dbo.History (DateTime, TagName, Value, QualityDetail) VALUES (@P1, @P2, @P3, 192)";
                insertCommand.CommandTimeout = 1000;

                var param1 = new SqlParameter("@P1", SqlDbType.DateTime);
                insertCommand.Parameters.Add(param1);

                var param2 = new SqlParameter("@P2", SqlDbType.NVarChar, 512);
                insertCommand.Parameters.Add(param2);

                var param3 = new SqlParameter("@P3", SqlDbType.Float);
                insertCommand.Parameters.Add(param3);

                insertCommand.Prepare();

                while (!ValueQueue.IsCompleted && ValueQueue.TryTake(out values, System.Threading.Timeout.Infinite))
                {
                    int retries = 0;

                    while (retries < 3)
                    {
                        insertCommand.Parameters["@P1"].Value = values.SampleDate.ToLocalTime();
                        insertCommand.Parameters["@P2"].Value = values.TagName;
                        insertCommand.Parameters["@P3"].Value = values.TagValue;

                        try
                        {
                            insertCommand.ExecuteNonQuery();
                            retries = 4;
                        }
                        catch (SqlException)
                        {
                            retries += 1;
                            sw.WriteLine("SQLException - Values: " + insertCommand.Parameters["@P1"].Value + ", " + insertCommand.Parameters["@P2"].Value + ", " + insertCommand.Parameters["@P3"].Value);
                        }
                    }
                }
            }
        }
    }

【问题讨论】:

  • SqlBulkCopy 类确实通过从内置的 SqlDataReader 读取进行了优化。与其将记录读取到并发队列然后将它们写回,是否可以仅在 @ 上使用 SqlBulkCopy 987654325@来自第一个函数?
  • 我试着走这条路。我的目标数据库实际上是专有的,并通过使用自定义 OLEDB 接口的链接服务器进行引用。批量插入不适用于链接服务器。
  • 两个函数的连接字符串是否都指向同一个 SQL 服务器(链接服务器是否与源数据库链接在同一台服务器上)?如果是这样,你为什么要通过 C#,为什么不直接在单个查询中进行查询?

标签: c# sql-server concurrent-queue


【解决方案1】:

您可以从正在检索的数据创建一个 DataTable 并使用 SqlBulkCopy 方法。请检查以下网址;

http://msdn.microsoft.com/en-us/library/ex21zs8x.aspx

【讨论】:

  • 我试着走这条路。我的目标数据库实际上是专有的,并通过使用自定义 OLEDB 接口的链接服务器进行引用。批量插入不适用于链接服务器。
  • 在 DatabaseHelper 类中使用 2 个不同的连接字符串怎么样?如果这不起作用,那么您也可以尝试 BulkInsert 方法:msdn.microsoft.com/en-us/library/ms175915.aspx
猜你喜欢
  • 2012-12-28
  • 1970-01-01
  • 2020-11-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-01-14
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多