【问题标题】:Create a multi threaded applications to run multiple queries in c#在 C# 中创建一个多线程应用程序以运行多个查询
【发布时间】:2017-02-07 22:51:26
【问题描述】:

我正在尝试构建一个异步运行查询的 Windows 窗体工具。 该应用程序有一个数据网格视图,可以运行 30 个可能的查询。用户检查他想要执行的查询,比如 10 个查询,然后点击一个按钮。 该应用程序有一个名为 maxthreads = 3 的变量(为了讨论),它指示可以使用多少线程来异步运行查询。查询在生产环境中运行,我们不希望同时运行太多线程使系统过载。每个查询平均运行 30 秒。 (有的 5 分钟,有的 2 秒) 在 datagridview 中有一个包含图标的图像列,该图标描述了每个查询的状态(0- 可以运行,1- 已选择运行,2- 正在运行,3- 成功完成,-1 错误) 每次查询开始和结束时,我都需要能够与 UI 进行通信。查询完成后,结果将显示在 Tabcontrol 中包含的 datagridview 中(每个查询一个选项卡)

方法:我正在考虑创建一些 maxthread 后台工作人员并让他们运行查询。当后台工作人员完成时,它会与 UI 通信并分配给一个新查询,依此类推,直到所有查询都运行完毕。

我尝试使用 assignmentWorker 将工作分派给后台工作人员,但不知道如何等待所有线程完成。一旦 bgw 完成,它会向 assignmentWorker 报告 RunWorkerCompleted 事件的进度,但该事件已经完成。

在 UI 线程中,我使用所有需要运行的查询调用分配工作者:

private void btnRunQueries_Click(object sender, EventArgs e)
    {
        if (AnyQueriesSelected())
        {
            tcResult.TabPages.Clear();

            foreach (DataGridViewRow dgr in dgvQueries.Rows)
            {
                if (Convert.ToBoolean(dgr.Cells["chk"].Value))
                {
                    Query q = new Query(dgr.Cells["ID"].Value.ToString(),
                        dgr.Cells["Name"].Value.ToString(),
                        dgr.Cells["FileName"].Value.ToString(),
                        dgr.Cells["ShortDescription"].Value.ToString(),
                        dgr.Cells["LongDescription"].Value.ToString(),
                        dgr.Cells["Level"].Value.ToString(),
                        dgr.Cells["Task"].Value.ToString(),
                        dgr.Cells["Importance"].Value.ToString(),
                        dgr.Cells["SkillSet"].Value.ToString(),
                        false,
                        new Dictionary<string, string>() 
                        { { "#ClntNb#", txtClntNum.Text }, { "#Staff#", "100300" } });

                    qryList.Add(q);
                }
            }
            assignmentWorker.RunWorkerAsync(qryList);
        }
        else
        {
            MessageBox.Show("Please select at least one query.",
                            "Warning",
                            MessageBoxButtons.OK,
                            MessageBoxIcon.Information);
        }
    }

这是 AssignmentWorker:

private void assignmentWorker_DoWork(object sender, DoWorkEventArgs e)
    {
        foreach (Query q in (List<Query>)e.Argument)
        {
            while (!q.Processed)
            {
                for (int threadNum = 0; threadNum < maxThreads; threadNum++)
                {
                    if (!threadArray[threadNum].IsBusy)
                    {
                        threadArray[threadNum].RunWorkerAsync(q);
                        q.Processed = true;
                        assignmentWorker.ReportProgress(1, q);
                        break;
                    }
                }

                //If all threads are being used, sleep awhile before checking again
                if (!q.Processed)
                {
                    Thread.Sleep(500);
                }
            }
        }
    }

所有 bgw 运行相同的事件:

private void backgroundWorkerFiles_DoWork(object sender, DoWorkEventArgs e)
    {
        try
        {
            Query qry = (Query)e.Argument;

            DataTable dtNew = DataAccess.RunQuery(qry).dtResult;

            if (dsQryResults.Tables.Contains(dtNew.TableName))
            {
                dsQryResults.Tables.Remove(dtNew.TableName);
            }

            dsQryResults.Tables.Add(dtNew);

            e.Result = qry;
        }
        catch (Exception ex)
        {

        }
    }

一旦 Query 返回并且 DataTable 已添加到数据集中:

private void backgroundWorkerFiles_RunWorkerCompleted(object sender, 
                                                    RunWorkerCompletedEventArgs e)
    {
        try
        {
            if (e.Error != null)
            {
                assignmentWorker.ReportProgress(-1, e.Result);
            }
            else
            {
                assignmentWorker.ReportProgress(2, e.Result);
            }
        }
        catch (Exception ex)
        {
            int o = 0;
        }
    }

我遇到的问题是分配工作人员在 bgw 完成之前完成,并且对 assignmentWorker.ReportProgress 的调用进入地狱(请原谅我的法语)。 如何等待所有已启动的 bgw 完成后再完成分配工作?

谢谢!

【问题讨论】:

  • 我不会将其拆分为后台工作人员和分配工作人员,这对于这项任务来说太复杂了。你可以有一个后台线程,queries to run 中的foreach 应该在 ThreadPool 上启动工作,或者等待工作线程数低于maxthreads - 并循环直到所有要处理的查询完成。要显示结果,后台任务应该使用Dispatcher 对 UI 进行适当的更新以开始在主 UI 线程上工作。

标签: c# multithreading backgroundworker


【解决方案1】:

正如the comment above 中所述,您的设计过于复杂。如果您有一个特定的最大数量的任务(查询)应该同时执行,您可以并且应该简单地创建该数量的工作人员,并让他们使用您的任务队列(或列表)中的任务,直到该队列为空。

缺少一个简洁明了地说明您的特定场景的良好Minimal, Complete, and Verifiable code example,因此提供可以直接解决您的问题的代码是不可行的。但是,这是一个使用 List&lt;T&gt; 的示例,就像您的原始代码一样,它将像我上面描述的那样工作:

using System;
using System.Collections.Generic;
using System.Threading.Tasks;

namespace TestSO42101517WaitAsyncTasks
{
    class Program
    {
        static void Main(string[] args)
        {
            Random random = new Random();
            int maxTasks = 30,
                maxActive = 3,
                maxDelayMs = 1000,
                currentDelay = -1;
            List<TimeSpan> taskDelays = new List<TimeSpan>(maxTasks);

            for (int i = 0; i < maxTasks; i++)
            {
                taskDelays.Add(TimeSpan.FromMilliseconds(random.Next(maxDelayMs)));
            }

            Task[] tasks = new Task[maxActive];
            object o = new object();

            for (int i = 0; i < maxActive; i++)
            {
                int workerIndex = i;

                tasks[i] = Task.Run(() =>
                {
                    DelayConsumer(ref currentDelay, taskDelays, o, workerIndex);
                });
            }

            Console.WriteLine("Waiting for consumer tasks");

            Task.WaitAll(tasks);

            Console.WriteLine("All consumer tasks completed");
        }

        private static void DelayConsumer(ref int currentDelay, List<TimeSpan> taskDelays, object o, int workerIndex)
        {
            Console.WriteLine($"worker #{workerIndex} starting");

            while (true)
            {
                TimeSpan delay;    
                int delayIndex;

                lock (o)
                {
                    delayIndex = ++currentDelay;
                    if (delayIndex < taskDelays.Count)
                    {
                        delay = taskDelays[delayIndex];
                    }
                    else
                    {
                        Console.WriteLine($"worker #{workerIndex} exiting");
                        return;
                    }
                }

                Console.WriteLine($"worker #{workerIndex} sleeping for {delay.TotalMilliseconds} ms, task #{delayIndex}");
                System.Threading.Thread.Sleep(delay);
            }
        }
    }
}

在您的情况下,每个工作人员都会向某个全局状态报告进度。您没有为您的“分配”工作人员显示ReportProgress 处理程序,所以我不能具体说明这会是什么样子。但大概它会涉及将-12 传递给某个知道如何处理这些值的方法(即,否则将是您的ReportProgress 处理程序)。

请注意,如果您对任务使用实际的队列数据结构,代码可以稍微简化一些,特别是在使用单个任务的情况下。这种方法看起来像这样:

using System;
using System.Collections.Concurrent;
using System.Threading.Tasks;

namespace TestSO42101517WaitAsyncTasks
{
    class Program
    {
        static void Main(string[] args)
        {
            Random random = new Random();
            int maxTasks = 30,
                maxActive = 3,
                maxDelayMs = 1000,
                currentDelay = -1;
            ConcurrentQueue<TimeSpan> taskDelays = new ConcurrentQueue<TimeSpan>();

            for (int i = 0; i < maxTasks; i++)
            {
                taskDelays.Enqueue(TimeSpan.FromMilliseconds(random.Next(maxDelayMs)));
            }

            Task[] tasks = new Task[maxActive];

            for (int i = 0; i < maxActive; i++)
            {
                int workerIndex = i;

                tasks[i] = Task.Run(() =>
                {
                    DelayConsumer(ref currentDelay, taskDelays, workerIndex);
                });
            }

            Console.WriteLine("Waiting for consumer tasks");

            Task.WaitAll(tasks);

            Console.WriteLine("All consumer tasks completed");
        }

        private static void DelayConsumer(ref int currentDelayIndex, ConcurrentQueue<TimeSpan> taskDelays, int workerIndex)
        {
            Console.WriteLine($"worker #{workerIndex} starting");

            while (true)
            {
                TimeSpan delay;

                if (!taskDelays.TryDequeue(out delay))
                {
                    Console.WriteLine($"worker #{workerIndex} exiting");
                    return;
                }

                int delayIndex = System.Threading.Interlocked.Increment(ref currentDelayIndex);

                Console.WriteLine($"worker #{workerIndex} sleeping for {delay.TotalMilliseconds} ms, task #{delayIndex}");
                System.Threading.Thread.Sleep(delay);
            }
        }
    }
}

【讨论】:

    猜你喜欢
    • 2021-12-23
    • 2018-07-03
    • 1970-01-01
    • 2020-05-27
    • 2018-08-13
    • 1970-01-01
    • 2020-12-13
    • 1970-01-01
    • 2013-08-20
    相关资源
    最近更新 更多