You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

.NET Semaphore任务退出时的高效清理方案咨询

问题描述

我有一段可配置并发消息处理数量的代码,使用SemaphoreSlim确保同时运行的任务数量不超过设定值。通常使用信号量时,我会捕获所有线程并确保最终全部完成,但此处任务数量可能随应用运行时长达到数千个。我希望在应用请求停止时,等待所有线程清理完成后再退出下方的ListenAsync方法。

public abstract class ContinuousListener
{
    private readonly ListenerOptions _listenerOptions;
    private readonly SemaphoreSlim _semaphore;

    public abstract Task ProcessMessageAsync(CancellationToken cancellationToken);

    public ContinuousListener(IOptions<ListenerOptions> listenerOptions)
    {
        _listenerOptions = listenerOptions.Value;
        _semaphore = new SemaphoreSlim(_listenerOptions.MaxConcurrentMessages);
    }

    public async Task ListenAsync(CancellationToken cancellationToken)
    {
        while (cancellationToken.IsCancellationRequested)
        {
           await _semaphore.WaitAsync(cancellationToken);

           var taskToWaitToCancel = ProcessSingleMessageAsync(cancellationToken);
        
           await Task.Delay(_listenerOptions.MessageDelay, cancellationToken);
       }
    }

    private async Task ProcessSingleMessageAsync(CancellationToken cancellationToken)
    {
        try
        {
            await ProcessMessageAsync(cancellationToken);
        }
        catch (Exception e)
        {
            // Log Exception and continue running
        }
        finally
        {
            _semaphore.Release();
        }
    }
}

请问针对taskToWaitToCancel是否有合适的任务清理模式?将任务加入列表会因大量已完成任务占用内存,我曾考虑编写仓储类跟踪当前任务,通过增删方法管理并调用取消操作以确保所有任务结束。

解决方案

首先纠正代码中的核心逻辑错误:原ListenAsync方法的循环条件while (cancellationToken.IsCancellationRequested)完全反转,应改为while (!cancellationToken.IsCancellationRequested),否则启动后会直接退出循环,无法处理任何消息。

针对任务跟踪与内存占用问题,推荐仅维护活跃任务集合的方案,既避免内存累积,又能实现优雅停机:

  1. 用线程安全集合跟踪未完成的活跃任务
  2. 任务完成后自动从集合中移除,确保集合只保留运行中的任务
  3. 收到取消信号时,先停止创建新任务,再等待所有活跃任务执行完毕

修改后的完整代码如下:

using System.Collections.Concurrent;

public abstract class ContinuousListener
{
    private readonly ListenerOptions _listenerOptions;
    private readonly SemaphoreSlim _semaphore;
    // 仅存储未完成任务的线程安全集合
    private readonly ConcurrentDictionary<Task, object> _activeTasks = new();

    public abstract Task ProcessMessageAsync(CancellationToken cancellationToken);

    public ContinuousListener(IOptions<ListenerOptions> listenerOptions)
    {
        _listenerOptions = listenerOptions.Value;
        _semaphore = new SemaphoreSlim(_listenerOptions.MaxConcurrentMessages);
    }

    public async Task ListenAsync(CancellationToken cancellationToken)
    {
        // 修正循环条件:未触发取消时持续处理消息
        while (!cancellationToken.IsCancellationRequested)
        {
            await _semaphore.WaitAsync(cancellationToken);

            var processingTask = ProcessSingleMessageAsync(cancellationToken);
            // 将任务加入活跃集合,并在完成后自动移除
            _activeTasks.TryAdd(processingTask, null);
            _ = processingTask.ContinueWith(task =>
            {
                _activeTasks.TryRemove(task, out _);
            }, TaskContinuationOptions.ExecuteSynchronously);
        
            await Task.Delay(_listenerOptions.MessageDelay, cancellationToken);
        }

        // 收到取消信号后,等待所有正在处理的任务完成
        if (_activeTasks.Count > 0)
        {
            await Task.WhenAll(_activeTasks.Keys);
        }
    }

    private async Task ProcessSingleMessageAsync(CancellationToken cancellationToken)
    {
        try
        {
            await ProcessMessageAsync(cancellationToken);
        }
        catch (Exception e)
        {
            // 保留日志逻辑
            // LogException(e);
        }
        finally
        {
            _semaphore.Release();
        }
    }
}

方案优势

  • 内存高效:集合仅保留未完成的活跃任务,已完成任务会自动被清理,不会累积占用内存
  • 线程安全:使用ConcurrentDictionary确保多线程环境下的集合操作安全
  • 简洁可靠:利用ContinueWith自动处理任务完成后的清理逻辑,无需手动维护增删操作
  • 优雅停机:收到取消信号后,先停止创建新任务,再等待所有正在处理的任务完成,保证消息处理的完整性

补充说明

  • 如果需要主动取消活跃任务,可以在停机逻辑中遍历_activeTasks.Keys,结合自定义取消令牌实现,但需确保ProcessMessageAsync方法能正确响应取消信号
  • 若担心ContinueWith的额外开销,也可以在ProcessSingleMessageAsync的finally块中移除任务,但需注意线程安全(异步方法的finally块可能在不同线程执行),因此仍推荐使用ConcurrentDictionary配合ContinueWith的方案

内容的提问来源于stack exchange,提问作者Nate

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 03:12:32