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

如何使这个异步方法调用工作?

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

    我尝试使用异步方法调用开发方法管道。管道的逻辑如下

    1. 集合中有n个数据必须输入管道中m个方法中
    2. 枚举t的集合
    3. 将第一个元素馈送给第一个方法
    4. 获取输出,异步将其馈送给第二个方法
    5. 同时,将集合的第二个元素馈送给第一个方法
    6. 完成第一个方法后,将结果发送给第二个方法(如果第二个方法仍在运行,则将结果放入其队列,并在第一个方法开始执行第三个元素)
    7. 当第二个方法完成执行时,从队列中取出第一个元素并执行,依此类推(每个方法都应该异步运行,没有人应该等待下一个方法完成)
    8. 在mth方法中,执行数据后,将结果存储到列表中
    9. 在mth方法中完成第n个元素后,将结果列表(n个结果)返回到第一个级别。

    我想出了一个如下的代码,但是它没有按预期工作,结果永远不会返回,而且它没有按应该的顺序执行。

    static class Program
        {
            static void Main(string[] args)
            {
                var list = new List<int> { 1, 2, 3, 4 };
                var result = list.ForEachPipeline(Add, Square, Add, Square);
                foreach (var element in result)
                {
                    Console.WriteLine(element);
                    Console.WriteLine("---------------------");
                }
                Console.ReadLine();
            }
    
            private static int Add(int j)
            {
                return j + 1;
            }
    
            private static int Square(int j)
            {
                return j * j;
            }
    
            internal static void AddNotify<T>(this List<T> list, T item)
            {
                Console.WriteLine("Adding {0} to the list", item);
                list.Add(item);
            }    
        }
    
        internal class Function<T>
        {
            private readonly Func<T, T> _func;
    
            private readonly List<T> _result = new List<T>();
            private readonly Queue<T> DataQueue = new Queue<T>();
            private bool _isBusy;
            static readonly object Sync = new object();
            readonly ManualResetEvent _waitHandle = new ManualResetEvent(false);
    
            internal Function(Func<T, T> func)
            {
                _func = func;
            }
    
            internal Function<T> Next { get; set; }
            internal Function<T> Start { get; set; }
            internal int Count;
    
            internal IEnumerable<T> Execute(IEnumerable<T> source)
            {
                var isSingle = true;
                foreach (var element in source) {
                    var result = _func(element);
                    if (Next != null)
                    {
                        Next.ExecuteAsync(result, _waitHandle);
                        isSingle = false;
                    }
                    else
                        _result.AddNotify(result);
                }
                if (!isSingle)
                    _waitHandle.WaitOne();
                return _result;
            }
    
    
            internal void ExecuteAsync(T element, ManualResetEvent resetEvent)
            {
                lock(Sync)
                {
                    if(_isBusy)
                    {
                        DataQueue.Enqueue(element);
                        return;
                    }
                    _isBusy = true;
    
                    _func.BeginInvoke(element, CallBack, resetEvent);
                }           
            }
    
            internal void CallBack(IAsyncResult result)
            {
                bool set = false;
                var worker = (Func<T, T>) ((AsyncResult) result).AsyncDelegate;
                var resultElement = worker.EndInvoke(result);
                var resetEvent = result.AsyncState as ManualResetEvent;
    
                lock(Sync)
                {
                    _isBusy = false;
                    if(Next != null)
                        Next.ExecuteAsync(resultElement, resetEvent);
                    else
                        Start._result.AddNotify(resultElement);
    
                    if(DataQueue.Count > 1)
                    {
                        var element = DataQueue.Dequeue();
                        ExecuteAsync(element, resetEvent);
                    }
                    if(Start._result.Count == Count)
                        set = true;
                }
                if(set)
                  resetEvent.Set();
            }
        }
    
        public static class Pipe
        {
            public static IEnumerable<T> ForEachPipeline<T>(this IEnumerable<T> source, params Func<T, T>[] pipes)
            {
                Function<T> start = null, previous = null;
                foreach (var function in pipes.Select(pipe => new Function<T>(pipe){ Count = source.Count()}))
                {
                    if (start == null)
                    {
                        start = previous = function;
                        start.Start = function;
                        continue;
                    }
                    function.Start = start;
                    previous.Next = function;
                    previous = function;
                }
                return start != null ? start.Execute(source) : null;
            }
        }
    

    你们能帮我把这件事做好吗?如果这种设计不适合实际的方法管道,请随意推荐不同的管道。

    编辑 :我必须严格遵守.NET 3.5。

    3 回复  |  直到 15 年前
        1
  •  1
  •   Wim Coenen    15 年前

    我没有立即在您的代码中发现问题,但您可能有点过于复杂了。这可能是一个简单的方法来做你想做的。

    public static class Pipe 
    {
       public static IEnumerable<T> Execute<T>(
          this IEnumerable<T> input, params Func<T, T>[] functions)
       {
          // each worker will put its result in this array
          var results = new T[input.Count()];
    
          // launch workers and return a WaitHandle for each one
          var waitHandles = input.Select(
             (element, index) =>
             {
                var waitHandle = new ManualResetEvent(false);
                ThreadPool.QueueUserWorkItem(
                   delegate
                   {
                      T result = element;
                      foreach (var function in functions)
                      {
                         result = function(result);
                      }
                      results[index] = result;
                      waitHandle.Set();
                   });
                return waitHandle;
             });
    
          // wait for each worker to finish
          foreach (var waitHandle in waitHandles)
          {
              waitHandle.WaitOne();
          }
          return results;
       }
    }
    

    这不会像在您自己的尝试中那样为管道的每个阶段创建锁。我忽略了这一点,因为它似乎没用。但是,您可以通过包装以下函数轻松地添加它:

    var wrappedFunctions = functions.Select(x => AddStageLock(x));
    

    哪里 AddStageLock 这是:

    private static Func<T,T> AddStageLock<T>(Func<T,T> function)
    {
       object stageLock = new object();
       Func<T, T> wrappedFunction =
          x =>
          {
             lock (stageLock)
             {
                return function(x);
             }
          };
       return wrappedFunction;
    }
    

    编辑: 这个 Execute 实现可能会比单线程执行慢,除非为每个单独的元素所做的工作使在线程池上创建等待句柄和调度任务的开销相形见绌,从而真正受益于多线程,您需要限制开销;.NET中的plinq4这个是在 partitioning the data

        2
  •  1
  •   VinayC    15 年前

    采用管道方法有什么特别的原因吗?在IMO中,为每个输入启动一个单独的线程,所有函数一个接一个地链接在一起,这样写起来更简单,执行起来也更快。例如,

    function T ExecPipe<T>(IEnumerable<Func<T, T>> pipe, T input)
    {
      T value = input;
      foreach(var f in pipe)
      {
        value = f(value);
      }
      return value;
    }
    
    var pipe = new List<Func<int, int>>() { Add, Square, Add, Square };
    var list = new List<int> { 1, 2, 3, 4 };
    foreach(var value in list)
    {
      ThreadPool.QueueUserWorkItem(o => ExecPipe(pipe, (int)o), value);
    }
    

    现在,说到您的代码,我相信对于使用m stage实现精确的管道,您必须有m个线程,因为每个阶段都可以并行执行-现在,一些线程可能是空闲的,因为I/P还没有到达它们。我不确定您的代码是否正在启动任何线程,以及在特定时间线程的计数。

        3
  •  0
  •   Slappy    15 年前

    为什么不为每次迭代断开一个线程,并将结果聚合到一个锁定资源中呢?你只需要做。可以使用plinq。 我认为你可能把方法误认为是资源。如果一个方法正在处理一个包含共享资源的关键块,那么您只需要锁定它。通过从中选择一个资源并从中切入一个新的线程,您就不再需要管理第二个方法了。

    即:方法x调用方法1,然后将值传递给方法2 ARR中的每个项目 异步(方法x(项));