【发布时间】:2017-08-07 16:48:29
【问题描述】:
我一直在构建一个服务,该服务使用Queue<string> 对象来处理文件以管理项目。
public partial class BasicQueueService : ServiceBase
{
private readonly EventWaitHandle completeHandle =
new EventWaitHandle(false, EventResetMode.ManualReset, "ThreadCompleters");
public BasicQueueService()
{
QueueManager = new Queue<string>();
}
public bool Stopping { get; set; }
private Queue<string> QueueManager { get; }
protected override void OnStart(string[] args)
{
Stopping = false;
ProcessFiles();
}
protected override void OnStop()
{
Stopping = true;
}
private void ProcessFiles()
{
while (!Stopping)
{
var count = QueueManager.Count;
for (var i = 0; i < count; i++)
{
//Check the Stopping Variable again.
if (Stopping) break;
var fileName = QueueManager.Dequeue();
if (string.IsNullOrWhiteSpace(fileName) || !File.Exists(fileName))
continue;
Console.WriteLine($"Processing {fileName}");
Task.Run(() =>
{
DoWork(fileName);
})
.ContinueWith(ThreadComplete);
}
if (Stopping) continue;
Console.WriteLine("Waiting for thread to finish, or 1 minute.");
completeHandle.WaitOne(new TimeSpan(0, 0, 15));
completeHandle.Reset();
}
}
partial void DoWork(string fileName);
private void ThreadComplete(Task task)
{
completeHandle.Set();
}
public void AddToQueue(string file)
{
//Called by FileWatcher/Manual classes, not included for brevity.
lock (QueueManager)
{
if (QueueManager.Contains(file)) return;
QueueManager.Enqueue(file);
}
}
}
在研究如何限制线程数量的同时(我尝试了一个手动类,其递增 int,但存在一个问题,即它在我的代码中没有正确递减),我遇到了 @987654321 @,这似乎更适合我想要实现的目标 - 具体来说,它允许我让框架处理线程/队列等。
现在这是我的服务:
public partial class BasicDataFlowService : ServiceBase
{
private readonly ActionBlock<string> workerBlock;
public BasicDataFlowService()
{
workerBlock = new ActionBlock<string>(file => DoWork(file), new ExecutionDataflowBlockOptions()
{
MaxDegreeOfParallelism = 32
});
}
public bool Stopping { get; set; }
protected override void OnStart(string[] args)
{
Stopping = false;
}
protected override void OnStop()
{
Stopping = true;
}
partial void DoWork(string fileName);
private void AddToDataFlow(string file)
{
workerBlock.Post(file);
}
}
这很好用。但是,我想确保一个文件只添加到TPL DataFlow 一次。使用Queue,我可以使用.Contains() 进行检查。有没有可以用于TPL DataFlow 的机制?
【问题讨论】:
-
无论是消费还是提交文件,都有责任不发布两次。如果您从目录中读取文件,您可以标记它们,或者按照@VMAtm 的建议缓存路径。但是,如果用户或其他客户正在提交它们,您需要将这些流程视为一项工作。其中每个文件代表具有单个结果的单个作业。
标签: c# multithreading tpl-dataflow dataflow