如何在C#中实现「即发即弃」异步FIFO队列?文档异步处理场景
解决异步文档处理队列的问题
嘿,我看了你用BlockingCollection实现异步任务队列时遇到的问题,咱们先拆解下当前代码里的核心坑:
- 每次添加任务就关闭队列:你在
Add方法里调用了_jobs.CompleteAdding(),这会直接把队列标记为“不再接受新任务”,后续再添加自然就不行了——这完全违背了“可持续添加任务”的需求。 - 重复启动处理逻辑:每次调用
Add都await ProcessQueue(),这会导致多个线程同时去消费同一个队列,任务会被乱序分割处理,还容易引发线程安全问题。 - 跨线程更新UI集合:
ObservableCollection不是线程安全的,直接在后台线程调用ExtractedDocuments.Add,UI绑定的时候肯定会炸(WPF/WinForms都不允许跨线程更新控件绑定的集合)。
下面给你两个简单易维护的方案,完全贴合你的小型应用需求:
方案一:用ConcurrentQueue实现按需启动的单任务处理
这个方案逻辑最直观,只有当队列有任务且没有正在处理的任务时,才启动处理循环,适合单任务串行处理的场景:
using System.Collections.Concurrent; using System.Collections.ObjectModel; using System.Threading; using System.Windows.Threading; public class QueueService { private readonly ConcurrentQueue<IDocument> _jobQueue = new ConcurrentQueue<IDocument>(); // 用原子操作标记是否正在处理,避免重复启动 private int _isProcessing = 0; // UI线程调度器,确保更新集合时在UI线程(WPF用Dispatcher,WinForms用SynchronizationContext.Current) private readonly Dispatcher _uiDispatcher; public ObservableCollection<IExtractedDocument> ExtractedDocuments { get; } public QueueService() { ExtractedDocuments = new ObservableCollection<IExtractedDocument>(); _uiDispatcher = Dispatcher.CurrentDispatcher; } public void Add(string filePath, List<Extra> extras) { var doc = new Document(filePath, extras); _jobQueue.Enqueue(doc); // 尝试启动处理,只有当前没在处理时才会成功 TryStartProcessing(); } private void TryStartProcessing() { // Interlocked.CompareExchange确保只有一个线程能启动处理逻辑 if (Interlocked.CompareExchange(ref _isProcessing, 1, 0) == 0) { // 用_ = 避免警告,这里不需要等待,让它后台运行 _ = ProcessQueueAsync(); } } private async Task ProcessQueueAsync() { try { // 循环处理队列里的所有任务 while (_jobQueue.TryDequeue(out var document)) { // 异步处理文档,完全不阻塞UI线程 var result = await service.ProcessDocument(document); // 必须在UI线程更新ObservableCollection,否则会抛出跨线程异常 await _uiDispatcher.InvokeAsync(() => { ExtractedDocuments.Add(result); }); Debug.WriteLine($"任务完成: {document.FilePath}"); } } catch (Exception ex) { Debug.WriteLine($"处理任务出错: {ex.Message}"); // 这里可以加日志、弹窗通知用户等错误处理 } finally { // 处理完成后重置标记,允许下次启动 Interlocked.Exchange(ref _isProcessing, 0); // 检查是否有新任务在处理过程中被添加进来,如果有就继续处理 if (!_jobQueue.IsEmpty) { TryStartProcessing(); } } } }
方案二:用BlockingCollection实现长期运行的处理线程
如果你更喜欢用BlockingCollection,可以改成这个版本——启动一个长期运行的后台线程,自动等待新任务,逻辑也很简洁:
using System.Collections.Concurrent; using System.Collections.ObjectModel; using System.Windows.Threading; public class QueueService { private readonly BlockingCollection<IDocument> _jobQueue = new BlockingCollection<IDocument>(); private Task _processingTask; private CancellationTokenSource _cts; private readonly Dispatcher _uiDispatcher; public ObservableCollection<IExtractedDocument> ExtractedDocuments { get; } public QueueService() { ExtractedDocuments = new ObservableCollection<IExtractedDocument>(); _uiDispatcher = Dispatcher.CurrentDispatcher; // 启动长期运行的处理任务 _cts = new CancellationTokenSource(); _processingTask = Task.Run(() => ProcessQueueAsync(_cts.Token)); } public void Add(string filePath, List<Extra> extras) { // 只有当队列还接受新任务时才添加 if (!_jobQueue.IsAddingCompleted) { var doc = new Document(filePath, extras); _jobQueue.Add(doc); } } private async Task ProcessQueueAsync(CancellationToken token) { try { // GetConsumingEnumerable会自动等待新任务,直到队列被标记为CompleteAdding foreach (var document in _jobQueue.GetConsumingEnumerable(token)) { var result = await service.ProcessDocument(document); await _uiDispatcher.InvokeAsync(() => ExtractedDocuments.Add(result)); Debug.WriteLine($"任务完成: {document.FilePath}"); } } catch (OperationCanceledException) { Debug.WriteLine("处理队列已被取消"); } catch (Exception ex) { Debug.WriteLine($"处理任务出错: {ex.Message}"); } } // 可选:当应用关闭时调用,停止处理队列 public async Task StopAsync() { _jobQueue.CompleteAdding(); _cts.Cancel(); await _processingTask; } }
扩展说明(任务状态)
如果后续要添加排队中/处理中/已完成的状态,只需要给任务实体加个状态属性就行:
public enum JobStatus { Queued, Processing, Completed, Failed } public class Job { public IDocument Document { get; set; } public JobStatus Status { get; set; } public string? ErrorMessage { get; set; } // 可选:添加进度属性 }
然后在入队时标记为Queued,开始处理时改成Processing,完成后改成Completed(失败则标记Failed),把ObservableCollection的泛型改成Job,UI就能直接绑定状态了。
并行处理的扩展(可选)
如果想要同时处理多个任务,只需要加个信号量控制并行数就行,比如最多同时处理2个:
// 在QueueService里添加 private readonly SemaphoreSlim _processingSemaphore = new SemaphoreSlim(2); // 修改处理逻辑为单个任务处理 private async Task ProcessSingleJobAsync(IDocument document) { await _processingSemaphore.WaitAsync(); try { var result = await service.ProcessDocument(document); await _uiDispatcher.InvokeAsync(() => ExtractedDocuments.Add(result)); } catch (Exception ex) { Debug.WriteLine($"处理任务出错: {ex.Message}"); } finally { _processingSemaphore.Release(); // 处理完一个任务后,尝试启动下一个 if (_jobQueue.TryDequeue(out var nextDoc)) { _ = ProcessSingleJobAsync(nextDoc); } } } // 修改Add方法里的启动逻辑 public void Add(string filePath, List<Extra> extras) { var doc = new Document(filePath, extras); _jobQueue.Enqueue(doc); // 如果有可用的并行槽位,就启动新的处理任务 if (_processingSemaphore.CurrentCount > 0) { if (_jobQueue.TryDequeue(out var nextDoc)) { _ = ProcessSingleJobAsync(nextDoc); } } }
这个并行方案同样简单,调整信号量的计数就能控制同时处理的任务数,完全不影响代码的可维护性。
内容的提问来源于stack exchange,提问作者Gil Sand
相关产品推荐
相关产品推荐

