场景描述:当同步遇到异步

在企业级开发中,经常会遇到这样一种”节奏不一致”的场景:客户端发起一个 HTTP 请求,期望同步得到结果;但后端的真实处理流程却是异步的——需要把任务投递到消息队列、等待第三方系统回调,或者由另一个工作进程处理后再通知。

举个典型例子:用户在前端点击”生成报表”,希望接口能立刻返回下载链接。但后端实际上是把任务丢给了后台 Worker,Worker 处理完成后会通过一个回调接口通知主服务。这时候主服务面临一个尴尬的局面:

  • 如果立即返回,用户拿不到结果;
  • 如果一直占着连接等待,又会导致请求堆积、线程池耗尽。

这就引出了一个经典的工程问题:如何把”同步等待”与”异步回调”桥接起来?

解决思路:挂起请求 - 等待回调 - 恢复响应

解决这个问题的核心思路其实很朴素,可以分成三步:

  1. 挂起请求:当请求到达时,不立即处理,而是为它生成一个唯一标识(比如 taskId),并创建一个”等待句柄”。请求线程进入等待状态,但不会阻塞 HTTP 连接所在的线程资源。
  2. 等待回调:把 taskId 发给异步处理方(队列、Worker、第三方)。当回调到来时,系统根据 taskId 找到对应的等待句柄。
  3. 恢复响应:通过句柄唤醒挂起的请求,把回调数据塞进响应里,原路返回给客户端。

在 .NET 中,实现这个”等待句柄”最优雅的工具就是 TaskCompletionSource<T>

基于 TaskCompletionSource 的实现

TaskCompletionSource<T>(简称 TCS)是 .NET 提供的一个用于手动控制 Task 状态的组件。它可以创建一个未完成的 Task,然后在任意时刻通过 SetResultSetExceptionSetCanceled 来完成它。这正好契合”挂起-唤醒”模型。

下面是一段最小化的核心实现:

// 用一个静态字典来保存所有挂起的请求
// key 是 taskId,value 是对应的 TaskCompletionSource
private static readonly ConcurrentDictionary<string, TaskCompletionSource<ReportResult>> _pendingRequests
    = new ConcurrentDictionary<string, TaskCompletionSource<ReportResult>>();

public async Task<IActionResult> GenerateReport(ReportRequest request)
{
    // 1. 生成唯一任务Id
    var taskId = Guid.NewGuid().ToString("N");

    // 2. 创建 TCS,并设置默认状态为 RunContinuationsAsynchronously
    //    这样可以避免回调线程直接同步执行 await 之后的代码,减少死锁风险
    var tcs = new TaskCompletionSource<ReportResult>(
        TaskCreationOptions.RunContinuationsAsynchronously);

    // 3. 把 tcs 登记到字典里,等待回调来唤醒
    _pendingRequests[taskId] = tcs;

    // 4. 把任务投递给后台 Worker(这里用消息队列示意)
    await _messageQueue.PublishAsync(new ReportTaskMessage
    {
        TaskId = taskId,
        Parameters = request.Parameters
    });

    // 5. 等待结果(可加超时),不阻塞当前线程
    var completedTask = await Task.WhenAny(tcs.Task, Task.Delay(TimeSpan.FromSeconds(30)));
    if (completedTask != tcs.Task)
    {
        // 超时,清理并返回超时提示
        _pendingRequests.TryRemove(taskId, out _);
        return StatusCode(408, new { message = "任务处理超时" });
    }

    // 6. 取出结果,从字典中移除
    _pendingRequests.TryRemove(taskId, out _);
    return Ok(tcs.Task.Result);
}

// 回调入口:后台 Worker 处理完成后调用这个接口
public IActionResult ReportCallback([FromBody] ReportResult result)
{
    // 根据 taskId 找到对应的 tcs,并设置结果,唤醒等待的请求
    if (_pendingRequests.TryGetValue(result.TaskId, out var tcs))
    {
        tcs.SetResult(result);
        return Ok();
    }
    // 找不到说明请求已超时被清理,记录日志即可
    return NotFound();
}

这段代码已经能跑通基本流程了:请求进来后被挂起,回调到来时被唤醒,结果原路返回。

超时处理:别让请求无限等待

上面的示例已经包含了简单的超时处理,但生产环境里更推荐用 CancellationTokenSource 来统一管理。它有两个好处:一是语义更清晰,二是可以和 ASP.NET Core 自身的请求取消机制联动。

public async Task<IActionResult> GenerateReport(ReportRequest request, CancellationToken ct)
{
    var taskId = Guid.NewGuid().ToString("N");

    // 创建带超时的 CTS,并和请求自身的取消令牌链接
    // 这样无论是超时还是客户端断开,都能统一处理
    using var cts = CancellationTokenSource.CreateLinkedTokenSource(ct);
    cts.CancelAfter(TimeSpan.FromSeconds(30));

    var tcs = new TaskCompletionSource<ReportResult>(
        TaskCreationOptions.RunContinuationsAsynchronously);

    // 注册取消回调:当取消发生时,让 tcs 进入取消状态
    using var registration = cts.Token.Register(() =>
    {
        // TrySetCanceled 避免重复设置(回调可能已经先到达)
        tcs.TrySetCanceled();
    });

    _pendingRequests[taskId] = tcs;

    try
    {
        await _messageQueue.PublishAsync(new ReportTaskMessage
        {
            TaskId = taskId,
            Parameters = request.Parameters
        });

        // 等待 tcs 完成,超时或取消都会抛 TaskCanceledException
        var result = await tcs.Task;
        return Ok(result);
    }
    catch (OperationCanceledException)
    {
        return StatusCode(408, new { message = "任务处理超时或被取消" });
    }
    finally
    {
        // 无论成功失败,都要清理字典,避免内存泄漏
        _pendingRequests.TryRemove(taskId, out _);
    }
}

关键点在于 finally 块里的清理:只要 tcs 还在字典里,它就一直占用内存。超时、异常、客户端断开都必须走清理逻辑,否则在高并发下会慢慢撑爆内存。

并发请求管理:用 ConcurrentDictionary 应对高并发

上面的实现用 ConcurrentDictionary 存储挂起的请求,这是处理并发的第一步。但在真正的生产环境里,还需要考虑几个细节:

1. 容量限制

如果回调方处理不过来,挂起的请求会无限堆积。可以加一个简单的容量检查:

// 限制同时挂起的请求数量,避免被压垮
private const int MaxPendingRequests = 1000;

if (_pendingRequests.Count >= MaxPendingRequests)
{
    return StatusCode(503, new { message = "系统繁忙,请稍后重试" });
}

2. 重复回调的容错

网络抖动可能导致同一个回调被投递多次。使用 TrySetResult 而不是 SetResult 就能避免”重复设置”的异常:

public IActionResult ReportCallback([FromBody] ReportResult result)
{
    if (_pendingRequests.TryGetValue(result.TaskId, out var tcs))
    {
        // TrySetResult 在 tcs 已经完成时返回 false,不会抛异常
        tcs.TrySetResult(result);
    }
    return Ok();
}

3. 集群环境下的状态共享

如果服务是集群部署,ConcurrentDictionary 只在单机有效。回调可能落到另一台机器上,找不到 tcs。这时候需要引入分布式缓存(如 Redis)或消息总线,把”哪个 taskId 在哪台机器上等待”这一映射存起来,再通过内部消息把回调路由到正确的节点。

完整代码示例:一个可用的桥接器

下面把前面的思路整合成一个独立、可复用的桥接器组件:

public interface IPendingRequestBridge<T>
{
    // 注册一个挂起的请求,返回 (taskId, 可等待的Task)
    (string TaskId, Task<T> Task) Register(TimeSpan timeout);

    // 完成指定任务,唤醒等待方
    bool Complete(string taskId, T result);

    // 让指定任务以异常方式完成
    bool Fail(string taskId, Exception ex);
}

public class PendingRequestBridge<T> : IPendingRequestBridge<T>
{
    // 单条挂起请求的上下文:包含 tcs 和清理用的取消注册
    private sealed class PendingEntry
    {
        public TaskCompletionSource<T> Tcs { get; } =
            new TaskCompletionSource<T>(TaskCreationOptions.RunContinuationsAsynchronously);
        public CancellationTokenSource Cts { get; } = new CancellationTokenSource();
    }

    private readonly ConcurrentDictionary<string, PendingEntry> _entries
        = new ConcurrentDictionary<string, PendingEntry>();

    public (string TaskId, Task<T> Task) Register(TimeSpan timeout)
    {
        var taskId = Guid.NewGuid().ToString("N");
        var entry = new PendingEntry();

        // 设置超时:超时后让 tcs 进入取消状态
        entry.Cts.CancelAfter(timeout);
        entry.Cts.Token.Register(() =>
        {
            // TrySet 避免与回调竞争
            entry.Tcs.TrySetCanceled();
            // 清理
            _entries.TryRemove(taskId, out _);
        });

        _entries[taskId] = entry;
        return (taskId, entry.Tcs.Task);
    }

    public bool Complete(string taskId, T result)
    {
        if (_entries.TryGetValue(taskId, out var entry))
        {
            // 成功完成时停止超时计时器
            entry.Cts.CancelAfter(TimeSpan.FromMilliseconds(-1));
            return entry.Tcs.TrySetResult(result);
        }
        return false;
    }

    public bool Fail(string taskId, Exception ex)
    {
        if (_entries.TryGetValue(taskId, out var entry))
        {
            return entry.Tcs.TrySetException(ex);
        }
        return false;
    }
}

在 Controller 里使用就非常清爽了:

public class ReportController : ControllerBase
{
    private readonly IPendingRequestBridge<ReportResult> _bridge;
    private readonly IMessageQueue _queue;

    public ReportController(IPendingRequestBridge<ReportResult> bridge, IMessageQueue queue)
    {
        _bridge = bridge;
        _queue = queue;
    }

    [HttpPost("report/generate")]
    public async Task<IActionResult> Generate(ReportRequest request)
    {
        // 注册挂起请求,超时30秒
        var (taskId, task) = _bridge.Register(TimeSpan.FromSeconds(30));

        // 投递给后台
        await _queue.PublishAsync(new ReportTaskMessage { TaskId = taskId });

        try
        {
            var result = await task;
            return Ok(result);
        }
        catch (OperationCanceledException)
        {
            return StatusCode(408, "超时");
        }
    }

    [HttpPost("report/callback")]
    public IActionResult Callback([FromBody] ReportResult result)
    {
        // 唤醒等待方
        _bridge.Complete(result.TaskId, result);
        return Ok();
    }
}

小结

把同步请求和异步回调桥接起来,本质上是借助 TaskCompletionSource 提供一个”手动可控的 Task”。核心要素有三个:

  • 挂起:用 TCS 创建一个未完成的 Task,并登记到字典里;
  • 唤醒:回调到来时通过 TrySetResult 完成 TCS;
  • 清理:超时、取消、异常都要确保从字典中移除,避免泄漏。

加上 CancellationTokenSource 做超时控制、ConcurrentDictionary 应对并发、容量限制防止过载,这套方案在真实项目中能稳定支撑每秒数千级的挂起请求。理解了这套模式,再去读类似 SignalR、gRPC 流式调用的实现,会发现它们底层也是同一套思路。