【问题标题】:Parallel.Foreach SQL querying sometimes results in ConnectionParallel.Foreach SQL 查询有时会导致连接
【发布时间】:2013-04-18 18:25:35
【问题描述】:

我需要加快在我的应用程序中执行 12 个查询。我从常规 foreach 切换到 Parallel.ForEach。但有时我会收到一条错误消息,提示“ExecuteReader 需要一个打开且可用的连接连接的当前状态为正在连接。”据我了解,由于 12 个查询中的许多查询都使用相同的 InitialCatalog,因此这 12 个查询中并没有真正的新连接,这可能是问题所在?我怎样才能解决这个问题? “sql”是“Sql”类型的列表——一个类只是一个字符串名称、字符串连接a和一个查询列表。代码如下:

 /// <summary>
    /// Connects to SQL, performs all queries and stores results in a list of DataTables
    /// </summary>
    /// <returns>List of data tables for each query in the config file</returns>
    public List<DataTable> GetAllData()
    {
        Stopwatch sw = new Stopwatch();
        sw.Start();
        List<DataTable> data = new List<DataTable>();

         List<Sql> sql=new List<Sql>();

        Sql one = new Sql();
         one.connection = "Data Source=XXX-SQL1;Initial Catalog=XXXDB;Integrated Security=True";
         one.name = "Col1";
         one.queries.Add("SELECT Name FROM [Reports]");
         one.queries.Add("SELECT Other FROM [Reports2]");
         sql.Add(one);

        Sql two = new Sql();
         two.connection = "Data Source=XXX-SQL1;Initial Catalog=XXXDB;Integrated Security=True";
         two.name = "Col2";
         two.queries.Add("SELECT AlternateName FROM [Reports1]");
         sql.Add(two);

         Sql three = new Sql();
         three.connection = "Data Source=YYY-SQL2;Initial Catalog=YYYDB;Integrated Security=True";
         three.name = "Col3";
         three.queries.Add("SELECT Frequency FROM Times");
         sql.Add(three);


        try
        {
            // ParallelOptions options = new ParallelOptions();
            //options.MaxDegreeOfParallelism = 3;
            // Parallel.ForEach(sql, options, s =>
            Parallel.ForEach(sql, s =>
            //foreach (Sql s in sql)
            {
                foreach (string q in s.queries)
                {
                    using (connection = new SqlConnection(s.connection))
                    {
                        connection.Open();
                        DataTable dt = new DataTable();
                        dt.TableName = s.name;
                        command = new SqlCommand(q, connection);
                        SqlDataAdapter adapter = new SqlDataAdapter();
                        adapter.SelectCommand = command;
                        adapter.Fill(dt);
                        //adapter.Dispose();

                        lock (data)
                        {
                            data.Add(dt);
                        }
                    }
                }
            }
            );
        }
        catch (Exception ex)
        {
            MessageBox.Show(ex.ToString(), "GetAllData error");
        }

        sw.Stop();
        MessageBox.Show(sw.Elapsed.ToString());

        return data;
    }

这是我创建的你需要的 Sql 类:

/// <summary>
/// Class defines a SQL connection and its respective queries
/// </summary>
public class Sql
{
    /// <summary>
    /// Name of the connection/query
    /// </summary>
    public string name { get; set; }
    /// <summary>
    /// SQL Connection string
    /// </summary>
    public string connection { get; set; }
    /// <summary>
    /// List of SQL queries for a connection
    /// </summary>
    public List<string> queries = new List<string>();
}

【问题讨论】:

    标签: c# sql parallel.foreach


    【解决方案1】:

    我会重构你的业务逻辑(连接到数据库)。

    public class SqlOperation
    {
        public SqlOperation()
        {
            Queries = new List<string>();
        }
    
        public string TableName { get; set; }
        public string ConnectionString { get; set; }
        public List<string> Queries { get; set; }
    }
    
    public static List<DataTable> GetAllData(IEnumerable<SqlOperation> sql)
    {
        var taskArray =
            sql.SelectMany(s =>
                s.Queries
                 .Select(query =>
                    Task.Run(() => //Task.Factory.StartNew for .NET 4.0
                        ExecuteQuery(s.ConnectionString, s.TableName, query))))
                .ToArray();
    
        try
        {
            Task.WaitAll(taskArray);
        }
        catch(AggregateException e)
        {
            MessageBox.Show(e.ToString(), "GetAllData error");
        }
    
        return taskArray.Where(t => !t.IsFaulted).Select(t => t.Result).ToList();
    }
    
    public static DataTable ExecuteQuery(string connectionString, string tableName, string query)
    {
        DataTable dataTable = null;
    
        using (var connection = new SqlConnection(connectionString))
        {
            dataTable = new DataTable();
            dataTable.TableName = tableName;
            using(var command = new SqlCommand(query, connection))
            {
                connection.Open();
    
                using(var adapter = new SqlDataAdapter())
                {
                    adapter.SelectCommand = command;
                    adapter.Fill(dataTable);
                }
            }
        }
    
         return dataTable;
    }
    

    【讨论】:

    • 你为什么要用Parallel.ForEach来启动一堆任务?您应该串行启动任务,或者让Parallel.ForEach 处理并行化。在这种情况下,我认为没有理由不做后者。
    • 我只是想用 Parallel.ForEach 替换 foreach 循环。基本上,由于其中一些查询需要一段时间,我想一次做不止 1 个查询。
    • @Romoku:看起来不错,但会产生错误:System.Threading.Tasks.Task 不包含“Result”的定义,并且没有扩展方法“Result”接受“System.Threading”类型的第一个参数可以找到 .Tasks.Task' 并且 'System.Threading.Tasks.Task' 不包含 'Run' 的定义。我需要什么不同的参考吗?
    • 对不起,昨天没来得及做出正确的回答,但我为你重构了它。您只需要引用System.Threading.Tasks
    • 改用Task.Factory.StartNew
    【解决方案2】:

    Ado.Net 有一个非常聪明的连接池,所以通常你应该只打开连接和关闭每个命令的连接,让池处理它们是否真的被打开或关闭。

    所以每个命令一个连接:

      Parallel.ForEach(sql, s=>
                //foreach (Sql s in sql)
                {
                    foreach (string q in s.queries)
                    {
                        using (connection = new SqlConnection(s.connection))
                        {
                            connection.Open();
                            DataTable dt = new DataTable();
                            dt.TableName = s.name;
                            command = new SqlCommand(q, connection);
                            SqlDataAdapter adapter = new SqlDataAdapter();
                            adapter.SelectCommand = command;
                            adapter.Fill(dt);
                            //adapter.Dispose();
    
                            lock(data){
                                data.Add(dt);
                            }
                        }
                    }
                }
    

    【讨论】:

    • 他还需要同步data.Add(dt);,因为List(T).Add不是线程安全的。
    • Hmmmmm...尝试一下(在使用之前将 foreach q 放入查询中会导致错误提示无法打开登录请求的数据库“xxx”。登录失败。用户“xxxxx 登录失败” “。每次。
    • @user1029770:抱歉,需要打开连接。见编辑。
    • 呵呵,现在我得到 ExecuteReader 需要一个打开且可用的连接。连接的当前状态每次都是连接错误。为什么我必须做 connection.Open() 当原来不需要的时候?对不起,我是新人!
    • 可能是因为数据适配器处理打开连接等。如果您提供更多代码,我将使其运行。
    【解决方案3】:

    您还可以在连接字符串中使用 MultipleActiveResultSets=true; 来支持多个阅读器

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-04-13
      • 2014-02-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-09-12
      • 1970-01-01
      • 2015-04-28
      相关资源
      最近更新 更多