【问题标题】:Implementing a buffer to write data from multiple threads?实现缓冲区以从多个线程写入数据?
【发布时间】:2011-06-08 02:35:33
【问题描述】:

我的程序使用一个迭代器来遍历一个映射,并产生一些工作线程来处理来自读取迭代器的,这一切都很好。现在,我想为每个点编写输出,为此我使用内存缓冲区来确保在将数据写入文件之前以正确的顺序从线程收集数据(通过另一个迭代器进行写入) :

public class MapMain
{
    // Multiple threads used here, each thread starts in Run() 
    // requests and processes map points

    public void Run()
    {
        // Get point from somewhere and process point
        int pointIndex = ...

        bufferWriter.StartPoint(pointIndex);

        // Perform a number of computations.
        // For simplicity, numberOfComputations = 1 in this example   
        bufferedWriter.BufferValue(pointIndex, value);

        bufferWriter.EndPoint(pointIndex); 
    }
}

我实现缓冲区的尝试:

public class BufferWriter
{
  private const int BufferSize = 4;

  private readonly IIterator iterator;
  private readonly float?[] bufferArray;
  private readonly bool[] bufferingCompleted;
  private readonly SortedDictionary<long, int> pointIndexToBufferIndexMap;
  private readonly object syncObject = new object();  

  private int bufferCount = 0;
  private int endBufferCount = 0;

  public BufferWriter(....)
  {
      iterator = ...
      bufferArray = new float?[BufferSize];
      bufferingCompleted = new bool[BufferSize];
      pointIndexToBufferIndexMap = new SortedDictionary<long, int>();
  }

  public void StartPoint(long pointIndex)
  {
    lock (syncObject)
    {
        if (bufferCount == BufferSize)
        {
            Monitor.Wait(syncObject);
        }

        pointIndexToBufferIndexMap.Add(pointIndex, bufferCount);   
        bufferCount++;
    }
  }

  public void BufferValue(long pointIndex, float value)
  {
      lock (syncObject)
      {
          int bufferIndex = pointIndexToBufferIndexMap[pointIndex];
          bufferArray[bufferIndex] = value;          
      }
  }

  public void EndPoint(long pointIndex)
  {
      lock (syncObject)
      {
          int bufferIndex = pointIndexToBufferIndexMap[pointIndex];
          bufferingCompleted[bufferIndex] = true;

          endBufferCount++;
          if (endBufferCount == BufferSize)
          {
              FlushBuffer();
              Monitor.PulseAll(syncObject);
          }
      }
  }

  private void FlushBuffer()
  {
      // Iterate in order of points
      foreach (long pointIndex in pointIndexToBufferIndexMap.Keys)
      {
          // Move iterator 
          iterator.MoveNext();

          int bufferIndex = pointIndexToBufferIndexMap[pointIndex];

          if (bufferArray[bufferIndex].HasValue)
          {                  
              iterator.Current = bufferArray[bufferIndex];

              // Clear to null
              bufferArray[bufferIndex] = null;                  
          }
      }

      bufferCount = 0;
      endBufferCount = 0;
      pointIndexToBufferIndexMap.Clear();
  }        
}

我正在寻找反馈来修复和纠正我的代码中的错误并解决任何性能问题:

[1] 简而言之:我有一个固定大小的缓冲区,它以某种随机顺序从多个线程处理点收集数据。当缓冲区完全被数据填满时,它必须被刷新。但是,如果我收集了点 0 到 9 但缺少点 8 怎么办?我的缓冲区已经满了,任何尝试使用缓冲区的点都会阻塞,直到执行刷新,这需要第 8 点。

[2] 缓冲区中值的顺序与值所引用的映射点的顺序不对应。如果是这种情况,那么我认为刷新会更容易(数组访问比 SortedDictionary 检索时间更快?)。此外,这可能允许我们重用刷新的槽来接收传入数据(循环缓冲区?)

但我想不出一个可行的模型来实现这一点。

[3] 缓冲区在刷新之前一直等到完全填满。在很多情况下,线程调用EndPoint() 并且iterator.Current 恰好引用了该点。 对于这一点,立即“写入”(即调用 'iterator.Current' 并枚举一次)可能更有意义,但如何做到这一点?

需要明确的是,BufferWriter 中的写入iterator 在其自身级别有一个缓冲区,用于缓存在写入输出之前对其Current 属性调用的值,但我不必担心。

我觉得整个事情需要从头开始重写!

任何帮助表示赞赏,谢谢。

【问题讨论】:

    标签: c# multithreading data-structures iterator buffer


    【解决方案1】:

    我不会“手动”进行并行处理,而是将其外包给 TPL 或 PLINQ。由于您在谈论地图,因此您有一组固定的点,您可以通过坐标枚举这些点,并让 PLINQ 担心并行性。

    示例:

    // first get your map points, could be just a lazy iterator over every map point
    IEnumerable<MapPoint> mapPoints = ...
    //Now use PLINQ to compute in parallel, maintain order
    var computedMapPoints = mapPoints.AsParallel()
                            .AsOrdered()
                            .Select(mappoint => ComputeMapPoint(mappoint)).ToList();
    

    【讨论】:

    • TPL 和 PLINQ 从 .NET 4.0 开始可用。还没有多少开发人员(由于各种原因)转向它。
    【解决方案2】:

    这应该是我的解决方案,虽然我还没有测试过。 添加新字段:

    private readonly Queue<AutoResetEvent> waitHandles = new Queue<AutoResetEvent>();
    

    两个 if(开始和结束)需要更改为:

    开始:

    if (bufferCount == BufferSize)
    {
        AutoResetEvent ev = new AutoResetEvent( false );
        waitHandles.Enqueue( ev );
        ev.WaitOne();
    }
    

    结束:

    if (endBufferCount == BufferSize)
    {
       FlushBuffer();
       for ( int i = 0; i < Math.Min( waitHandles.Count, BufferSize ); ++i )
       {
          waitHandles.Dequeue().Set();
       }
    }
    

    【讨论】:

    • 我当前的实现仅使用syncObj 限制对缓冲区的访问。如果一个线程试图通过调用 StartPoint() 来使用缓冲区,并且它恰好是满的,那么线程将不得不等待直到缓冲区被刷新。
    • 简而言之:我有一个固定大小的缓冲区,它以某种随机顺序从多个线程处理点收集数据。当缓冲区完全填满数据时,它必须被刷新,对吗?但是,如果我收集了点 0 到 9 但缺少点 8 怎么办?我的缓冲区已经满了,任何尝试使用缓冲区的点都会阻塞,直到执行刷新,这需要第 8 点
    • 我可以看到代码中的一些缺陷。 1.如果要邀请另外4个线程计算,只需要邀请4个线程,但PulseAll释放所有等待线程(即1000)。释放时,它们使用相同的 bufferCount 值 == 0(在 FlushBuffer 方法中归零):pointIndexToBufferIndexMap.Add(pointIndex, bufferCount);//bufferCount == 0 for ALL threads bufferCount++;//这是我认为的混乱编辑:不要不知道如何正确格式化评论。对不起。
    • 你检查我更新的解决方案了吗?我没有收到您的任何反馈。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-04-02
    • 1970-01-01
    相关资源
    最近更新 更多