ASP.NET Core Web API长时任务如何实现客户端状态通知与查询?
实现可跟踪状态的异步耗时任务方案
针对ASP.NET Core 6 Web API中耗时任务的状态跟踪需求,以下是具体实现方案:
1. 定义核心模型
先统一任务相关的数据结构:
// 任务请求模型(对应原PushInput) public class TaskRequest { // 替换为你的实际请求字段 public string InputData { get; set; } = string.Empty; } // 任务状态枚举 public enum TaskStatus { Pending, Running, Completed, Failed } // 任务信息模型,用于跟踪状态与结果 public class TaskInfo { public Guid TaskId { get; set; } public TaskStatus Status { get; set; } public TaskResult? Result { get; set; } public DateTime CreatedAt { get; set; } public string? ErrorMessage { get; set; } } // 任务结果模型(对应原PushOutput) public class TaskResult { // 替换为你的实际结果字段 public string OutputData { get; set; } = string.Empty; }
2. 实现任务状态管理器
使用线程安全的内存存储管理任务状态(新手友好,生产环境可替换为数据库/Redis持久化):
public interface ITaskManager { Guid SubmitTask(TaskRequest request); TaskInfo? GetTaskStatus(Guid taskId); void UpdateTaskStatus(Guid taskId, TaskStatus status); void SetTaskResult(Guid taskId, TaskResult result); void SetTaskFailure(Guid taskId, string errorMessage); } public class InMemoryTaskManager : ITaskManager { private readonly ConcurrentDictionary<Guid, TaskInfo> _tasks = new(); public Guid SubmitTask(TaskRequest request) { var taskId = Guid.NewGuid(); _tasks.TryAdd(taskId, new TaskInfo { TaskId = taskId, Status = TaskStatus.Pending, CreatedAt = DateTime.UtcNow }); return taskId; } public TaskInfo? GetTaskStatus(Guid taskId) { _tasks.TryGetValue(taskId, out var taskInfo); return taskInfo; } public void UpdateTaskStatus(Guid taskId, TaskStatus status) { if (_tasks.TryGetValue(taskId, out var taskInfo)) { taskInfo.Status = status; } } public void SetTaskResult(Guid taskId, TaskResult result) { if (_tasks.TryGetValue(taskId, out var taskInfo)) { taskInfo.Status = TaskStatus.Completed; taskInfo.Result = result; } } public void SetTaskFailure(Guid taskId, string errorMessage) { if (_tasks.TryGetValue(taskId, out var taskInfo)) { taskInfo.Status = TaskStatus.Failed; taskInfo.ErrorMessage = errorMessage; } } }
3. 实现后台任务执行服务
基于BackgroundService和通道队列处理异步任务,避免阻塞API请求:
public class TaskProcessingService : BackgroundService { private readonly Channel<(Guid TaskId, TaskRequest Request)> _taskChannel; private readonly ITaskManager _taskManager; private readonly ILogger<TaskProcessingService> _logger; public TaskProcessingService(ITaskManager taskManager, ILogger<TaskProcessingService> logger) { _taskManager = taskManager; _logger = logger; // 创建有界通道,避免内存溢出,容量可根据业务调整 _taskChannel = Channel.CreateBounded<(Guid, TaskRequest)>(new BoundedChannelOptions(100) { FullMode = BoundedChannelFullMode.Wait }); } // 供API调用,提交任务到队列 public async ValueTask SubmitTaskAsync(Guid taskId, TaskRequest request) { await _taskChannel.Writer.WriteAsync((taskId, request)); } protected override async Task ExecuteAsync(CancellationToken stoppingToken) { await foreach (var (taskId, request) in _taskChannel.Reader.ReadAllAsync(stoppingToken)) { try { _taskManager.UpdateTaskStatus(taskId, TaskStatus.Running); _logger.LogInformation("开始执行任务:{TaskId}", taskId); // 替换为你的实际耗时业务逻辑 await Task.Delay(50000, stoppingToken); // 生成任务结果,替换为你的实际结果生成逻辑 var result = new TaskResult { OutputData = $"处理完成:{request.InputData}" }; _taskManager.SetTaskResult(taskId, result); _logger.LogInformation("任务执行完成:{TaskId}", taskId); } catch (OperationCanceledException) { _taskManager.SetTaskFailure(taskId, "任务被取消"); _logger.LogWarning("任务被取消:{TaskId}", taskId); } catch (Exception ex) { _taskManager.SetTaskFailure(taskId, ex.Message); _logger.LogError(ex, "任务执行失败:{TaskId}", taskId); } } } }
4. 编写API控制器
提供任务提交和状态查询接口:
[ApiController] [Route("api/[controller]")] public class TasksController : ControllerBase { private readonly ITaskManager _taskManager; private readonly TaskProcessingService _taskProcessingService; public TasksController(ITaskManager taskManager, TaskProcessingService taskProcessingService) { _taskManager = taskManager; _taskProcessingService = taskProcessingService; } /// <summary> /// 提交耗时任务 /// </summary> [HttpPost] public IActionResult SubmitTask([FromBody] TaskRequest request) { var taskId = _taskManager.SubmitTask(request); // 异步提交到后台服务,不阻塞当前请求 _ = _taskProcessingService.SubmitTaskAsync(taskId, request); return AcceptedAtAction(nameof(GetTaskStatus), new { taskId }, new { TaskId = taskId }); } /// <summary> /// 查询任务状态与结果 /// </summary> [HttpGet("{taskId:guid}")] public IActionResult GetTaskStatus(Guid taskId) { var taskInfo = _taskManager.GetTaskStatus(taskId); if (taskInfo == null) { return NotFound(); } return Ok(taskInfo); } }
5. 注册服务到DI容器
在Program.cs中添加服务注册:
var builder = WebApplication.CreateBuilder(args); // 添加控制器服务 builder.Services.AddControllers(); // 注册任务管理器和后台任务服务 builder.Services.AddSingleton<ITaskManager, InMemoryTaskManager>(); builder.Services.AddHostedService<TaskProcessingService>(); var app = builder.Build(); // 配置中间件 app.UseHttpsRedirection(); app.UseAuthorization(); app.MapControllers(); app.Run();
扩展方向
- 持久化:将内存存储替换为SQL Server、Redis等,避免API重启丢失任务状态
- 任务取消:新增取消接口,在后台服务中跟踪任务取消信号
- 超时控制:为耗时任务添加超时逻辑,避免无限期运行
- 重试机制:针对失败任务添加自动重试逻辑,配置重试次数与间隔
内容的提问来源于stack exchange,提问作者BeingPunekar
相关产品推荐
相关产品推荐

