场景描述:当同步遇到异步
在企业级开发中,经常会遇到这样一种”节奏不一致”的场景:客户端发起一个 HTTP 请求,期望同步得到结果;但后端的真实处理流程却是异步的——需要把任务投递到消息队列、等待第三方系统回调,或者由另一个工作进程处理后再通知。
举个典型例子:用户在前端点击”生成报表”,希望接口能立刻返回下载链接。但后端实际上是把任务丢给了后台 Worker,Worker 处理完成后会通过一个回调接口通知主服务。这时候主服务面临一个尴尬的局面:
- 如果立即返回,用户拿不到结果;
- 如果一直占着连接等待,又会导致请求堆积、线程池耗尽。
这就引出了一个经典的工程问题:如何把”同步等待”与”异步回调”桥接起来?
解决思路:挂起请求 - 等待回调 - 恢复响应
解决这个问题的核心思路其实很朴素,可以分成三步:
- 挂起请求:当请求到达时,不立即处理,而是为它生成一个唯一标识(比如
taskId),并创建一个”等待句柄”。请求线程进入等待状态,但不会阻塞 HTTP 连接所在的线程资源。 - 等待回调:把
taskId发给异步处理方(队列、Worker、第三方)。当回调到来时,系统根据taskId找到对应的等待句柄。 - 恢复响应:通过句柄唤醒挂起的请求,把回调数据塞进响应里,原路返回给客户端。
在 .NET 中,实现这个”等待句柄”最优雅的工具就是 TaskCompletionSource<T>。
基于 TaskCompletionSource 的实现
TaskCompletionSource<T>(简称 TCS)是 .NET 提供的一个用于手动控制 Task 状态的组件。它可以创建一个未完成的 Task,然后在任意时刻通过 SetResult、SetException、SetCanceled 来完成它。这正好契合”挂起-唤醒”模型。
下面是一段最小化的核心实现:
// 用一个静态字典来保存所有挂起的请求
// 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 流式调用的实现,会发现它们底层也是同一套思路。