【发布时间】: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