ASP.NET中CQRS与外部进程托管长时任务的通知需求
针对你这个长任务导出后的在线用户通知需求,我整理了几个贴合你现有架构的实用方案,你可以根据自己的技术栈和流量情况选择:
方案1:WebSocket(推荐,实时性最优)
WebSocket的双向通信特性完美适配这种后端主动推通知的场景,能保证用户第一时间收到任务完成消息。
实现步骤:
- 用户打开页面时,前端立即和主Web应用建立WebSocket连接,同时带上用户标识(比如用户ID、会话ID)和任务ID(触发导出时前端可获取到)。
- 外部自托管WebAPI完成导出任务后,调用主Web应用的一个回调接口,传递任务ID、用户ID和任务完成状态。
- 主Web应用收到回调后,通过WebSocket连接池找到对应在线用户的连接,推送通知消息。
代码示例(以ASP.NET Core+JS为例):
前端:// 页面加载时建立连接,假设已拿到taskId和当前用户ID const socket = new WebSocket(`ws://your-domain/ws/task-notify?taskId=${taskId}&userId=${currentUserId}`); socket.onmessage = (event) => { const msg = JSON.parse(event.data); if(msg.type === 'task-ready') { // 自定义通知逻辑:比如弹层提示、更新页面按钮状态 alert(`导出任务「${msg.taskName}」已完成,可前往下载!`); } }; // 页面关闭时主动断开连接 window.addEventListener('beforeunload', () => socket.close());后端:
// WebSocket端点,用于维护用户连接 [Route("ws/task-notify")] public async Task HandleWebSocket(HttpContext context) { if (!context.WebSockets.IsWebSocketRequest) { context.Response.StatusCode = 400; return; } var webSocket = await context.WebSockets.AcceptWebSocketAsync(); var taskId = context.Request.Query["taskId"]; var userId = context.Request.Query["userId"]; var connectionKey = $"{userId}_{taskId}"; // 将连接存入全局连接池(可用ConcurrentDictionary) _activeWebSocketConnections.TryAdd(connectionKey, webSocket); try { // 保持连接存活,监听关闭信号 var buffer = new byte[1024 * 4]; WebSocketReceiveResult result; do { result = await webSocket.ReceiveAsync(new ArraySegment<byte>(buffer), CancellationToken.None); } while (!result.CloseStatus.HasValue); await webSocket.CloseAsync(result.CloseStatus.Value, result.CloseStatusDescription, CancellationToken.None); } finally { // 连接关闭后从池移除 _activeWebSocketConnections.TryRemove(connectionKey, out _); } } // 供外部WebAPI调用的回调接口 [HttpPost("task/callback/complete")] public async Task<IActionResult> TaskCompleteCallback([FromBody] TaskCompletionDto dto) { var connectionKey = $"{dto.UserId}_{dto.TaskId}"; if (_activeWebSocketConnections.TryGetValue(connectionKey, out var webSocket)) { var notifyMsg = JsonSerializer.Serialize(new { type = "task-ready", taskName = dto.TaskName, downloadUrl = dto.DownloadUrl }); var buffer = Encoding.UTF8.GetBytes(notifyMsg); await webSocket.SendAsync(new ArraySegment<byte>(buffer), WebSocketMessageType.Text, true, CancellationToken.None); } return Ok(); } // 辅助DTO public class TaskCompletionDto { public string TaskId { get; set; } public string UserId { get; set; } public string TaskName { get; set; } public string DownloadUrl { get; set; } }优缺点:
✅ 实时性拉满,无延迟;后端主动推送,无无效请求
❌ 需要维护连接池,需处理断连重连(可添加心跳机制)
方案2:前端轮询(简单易实现,适合小流量场景)
如果不想引入WebSocket这类复杂机制,轮询是最直接的方案,前端定期查询任务状态。
实现步骤:
- 用户触发导出后,前端拿到任务ID,启动定时器每隔N秒请求主应用的任务状态接口。
- 外部WebAPI完成任务后,更新数据库/缓存中的任务状态为「已完成」,并写入下载链接。
- 前端查询到「已完成」状态后,显示通知并停止轮询。
代码示例:
前端:const taskId = "xxx"; // 触发导出时获取的任务ID const pollTimer = setInterval(async () => { try { const res = await fetch(`/api/task/status?taskId=${taskId}`); const status = await res.json(); if(status.isCompleted) { clearInterval(pollTimer); // 展示通知 document.getElementById('notify').innerHTML = ` <div class="alert alert-success"> 导出任务已完成!<a href="${status.downloadUrl}" target="_blank">点击下载</a> </div> `; } } catch(err) { console.error('查询任务状态失败:', err); } }, 10000); // 每10秒查询一次 // 页面关闭时清理定时器 window.addEventListener('beforeunload', () => clearInterval(pollTimer));后端:
[HttpGet("api/task/status")] public IActionResult GetTaskStatus(string taskId) { // 从数据库/缓存读取任务状态 var task = _taskRepository.GetById(taskId); if(task == null) return NotFound(); return Ok(new { isCompleted = task.IsCompleted, downloadUrl = task.IsCompleted ? task.DownloadUrl : null }); }优缺点:
✅ 实现零复杂度,无需额外技术依赖
❌ 存在无效请求,实时性取决于轮询间隔(间隔短增加服务器压力,间隔长通知延迟)
方案3:Server-Sent Events(SSE)
SSE是单向的服务器推送技术,比WebSocket简单,适合只需要后端推消息的场景。
实现步骤:
- 前端建立SSE连接,带上任务ID和用户ID。
- 后端保持连接打开,定时检查任务状态,一旦完成就通过连接发送通知。
- 前端收到通知后展示并关闭连接。
代码示例:
前端:const eventSource = new EventSource(`/api/task/sse-notify?taskId=${taskId}&userId=${currentUserId}`); eventSource.addEventListener('task-completed', (event) => { const data = JSON.parse(event.data); alert(`导出任务「${data.taskName}」已完成!`); eventSource.close(); // 完成后关闭连接 }); // 处理连接错误 eventSource.onerror = () => { console.error('SSE连接异常,正在重连...'); eventSource.close(); // 可添加重连逻辑 };后端:
[HttpGet("api/task/sse-notify")] public async Task SendTaskNotification(HttpContext context) { context.Response.Headers.Add("Content-Type", "text/event-stream"); context.Response.Headers.Add("Cache-Control", "no-cache"); context.Response.Headers.Add("Connection", "keep-alive"); var taskId = context.Request.Query["taskId"]; var cancellationToken = context.RequestAborted; while (!cancellationToken.IsCancellationRequested) { var task = _taskRepository.GetById(taskId); if(task != null && task.IsCompleted) { var message = $"event: task-completed\ndata: {JsonSerializer.Serialize(new { taskName = task.Name })}\n\n"; await context.Response.WriteAsync(message); await context.Response.Body.FlushAsync(); break; } // 每隔5秒检查一次 await Task.Delay(5000, cancellationToken); } }优缺点:
✅ 比WebSocket简单,无需处理双向通信
❌ 仅支持后端推消息,部分浏览器有连接数量限制
额外注意事项
- 任务状态持久化:无论用哪种方案,都要将任务状态(完成状态、下载链接等)存入数据库或缓存,确保外部API完成任务后能正确更新状态。
- 回调可靠性:外部WebAPI调用主应用回调接口时,要添加重试机制(比如用消息队列兜底),避免因网络波动导致通知丢失。
- 用户匹配准确性:建议用「用户ID+会话ID」作为连接标识,避免同一用户多端登录时通知错误。
- 异常处理:前端要处理断连、网络错误等情况,后端要清理无效连接,避免资源泄漏。
内容的提问来源于stack exchange,提问作者Coffka
相关产品推荐
相关产品推荐

