代码之家  ›  专栏  ›  技术社区  ›  Anindya Chatterjee

如何并行运行一组函数并在完成时等待结果?

  •  4
  • Anindya Chatterjee  · 技术社区  · 15 年前

    我需要在sametime异步运行一组重函数,并将结果填充到列表中。这是伪代码:

    List<TResult> results = new List<TResults>();
    List<Func<T, TResult>> tasks = PopulateTasks();
    
    foreach(var task in tasks)
    {
        // Run Logic in question
        1. Run each task asynchronously/parallely
        2. Put the results in the results list upon each task completion
    }
    
    Console.WriteLine("All tasks completed and results populated");
    

    我需要内在的逻辑 foreach 博克。你们能帮帮我吗?

    提前谢谢。

    6 回复  |  直到 15 年前
        1
  •  4
  •   Ohad Schneider    15 年前

    一个简单的3.5实现可以如下所示

    List<TResult> results = new List<TResults>();
    List<Func<T, TResult>> tasks = PopulateTasks();
    
    ManualResetEvent waitHandle = new ManualResetEvent(false);
    void RunTasks()
    {
        int i = 0;
        foreach(var task in tasks)
        {
            int captured = i++;
            ThreadPool.QueueUserWorkItem(state => RunTask(task, captured))
        }
    
        waitHandle.WaitOne();
    
        Console.WriteLine("All tasks completed and results populated");
    }
    
    private int counter;
    private readonly object listLock = new object();
    void RunTask(Func<T, TResult> task, int index)
    {
        var res = task(...); //You haven't specified where the parameter comes from
        lock (listLock )
        {
           results[index] = res;
        }
        if (InterLocked.Increment(ref counter) == tasks.Count)
            waitHandle.Set();
    }
    
        2
  •  4
  •   wRAR    15 年前
    List<Func<T, TResult>> tasks = PopulateTasks();
    TResult[] results = new TResult[tasks.Length];
    Parallel.For(0, tasks.Count, i =>
        {
            results[i] = tasks[i]();
        });
    

    TPL for 3.5 apparently exists .

        3
  •  1
  •   Adrian Zanescu    15 年前
        public static IList<IAsyncResult> RunAsync<T>(IEnumerable<Func<T>> tasks)
        {
            List<IAsyncResult> asyncContext = new List<IAsyncResult>();
            foreach (var task in tasks)
            {
                asyncContext.Add(task.BeginInvoke(null, null));
            }
            return asyncContext;
        }
    
        public static IEnumerable<T> WaitForAll<T>(IEnumerable<Func<T>> tasks, IEnumerable<IAsyncResult> asyncContext)
        {
            IEnumerator<IAsyncResult> iterator = asyncContext.GetEnumerator();
            foreach (var task in tasks)
            {
                iterator.MoveNext();
                yield return task.EndInvoke(iterator.Current);
            }
        }
    
        public static void Main()
        {
            var tasks = Enumerable.Repeat<Func<int>>(() => ComputeValue(), 10).ToList();
    
            var asyncContext = RunAsync(tasks);
            var results = WaitForAll(tasks, asyncContext);
            foreach (var result in results)
            {
                Console.WriteLine(result);
            }
        }
    
        public static int ComputeValue()
        {
            Thread.Sleep(1000);
            return Guid.NewGuid().ToByteArray().Sum(a => (int)a); 
        }
    
        4
  •  1
  •   Adrian Zanescu    15 年前

    另一个变体是一个小型的未来模式实现:

        public class Future<T>
        {
            public Future(Func<T> task)
            {
                Task = task;
                _asyncContext = task.BeginInvoke(null, null);
            }
    
            private IAsyncResult _asyncContext;
    
            public Func<T> Task { get; private set; }
            public T Result
            {
                get
                {
                    return Task.EndInvoke(_asyncContext);
                }
            }
    
            public bool IsCompleted
            {
                get { return _asyncContext.IsCompleted; }
            }
        }
    
        public static IList<Future<T>> RunAsync<T>(IEnumerable<Func<T>> tasks)
        {
            List<Future<T>> asyncContext = new List<Future<T>>();
            foreach (var task in tasks)
            {
                asyncContext.Add(new Future<T>(task));
            }
            return asyncContext;
        }
    
        public static IEnumerable<T> WaitForAll<T>(IEnumerable<Future<T>> futures)
        {
            foreach (var future in futures)
            {
                yield return future.Result;
            }
        }
    
        public static void Main()
        {
            var tasks = Enumerable.Repeat<Func<int>>(() => ComputeValue(), 10).ToList();
    
            var futures = RunAsync(tasks);
            var results = WaitForAll(futures);
            foreach (var result in results)
            {
                Console.WriteLine(result);
            }
        }
    
        public static int ComputeValue()
        {
            Thread.Sleep(1000);
            return Guid.NewGuid().ToByteArray().Sum(a => (int)a);
        }
    
        5
  •  0
  •   gbjbaanb    15 年前

        6
  •  0
  •   Ed Power    15 年前

    在单独的工作实例中进行处理,每个工作实例都在各自的线程上。使用回调返回结果并向调用进程发送线程已完成的信号。使用字典跟踪正在运行的线程。如果有很多线程,则应加载队列并在旧线程完成时启动新线程。在本例中,所有线程都是在启动任何线程之前创建的,以防止在启动最终线程之前运行的线程计数降至零的竞争条件。

        Dictionary<int, Thread> activeThreads = new Dictionary<int, Thread>();
        void LaunchWorkers()
        {
            foreach (var task in tasks)
            {
                Worker worker = new Worker(task, new WorkerDoneDelegate(ProcessResult));
                Thread thread = new Thread(worker.Done);
                thread.IsBackground = true;
                activeThreads.Add(thread.ManagedThreadId, thread);
            }
            lock (activeThreads)
            {
                activeThreads.Values.ToList().ForEach(n => n.Start());
            }
        }
    
        void ProcessResult(int threadId, TResult result)
        {
            lock (results)
            {
                results.Add(result);
            }
            lock (activeThreads)
            {
                activeThreads.Remove(threadId);
                // done when activeThreads.Count == 0
            }
        }
    }
    
    public delegate void WorkerDoneDelegate(object results);
    class Worker
    {
        public WorkerDoneDelegate Done;
        public void Work(Task task, WorkerDoneDelegate Done)
        {
            // process task
            Done(Thread.CurrentThread.ManagedThreadId, result);
        }
    }