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

基于C# .NET Core MVC的大数据量后台任务处理及用户通知问询

完全可以实现这个需求!我之前在做.NET Core MVC项目时也遇到过类似的批量数据处理场景,下面给你拆解具体的实现思路和关键步骤:

核心思路概述

核心逻辑是把长时间运行的任务从同步HTTP请求中剥离:先判断任务预估/实际耗时是否超过阈值,若超过则将任务放入后台队列异步执行,前端立即返回任务状态提示;后台任务执行过程中跟踪状态,完成后通过邮件、站内信或实时推送通知用户结果。

具体实现步骤

1. 任务耗时判断

有两种常见方式来判断是否需要转入后台:

  • 提前预估:根据数据量计算大致耗时(比如每处理1000条数据库记录需要1秒,超过30秒就转后台)
  • 实时监控:用Stopwatch在同步执行过程中监控耗时,一旦超过阈值就中断当前操作,转后台执行(注意要保证数据一致性,比如用事务回滚已执行的部分,或设计幂等性逻辑避免重复处理)

2. 实现后台任务队列

.NET Core原生提供了BackgroundService来实现后台任务,结合自定义任务队列可以灵活管理异步任务:

第一步:定义任务队列接口与实现

public interface IBackgroundTaskQueue
{
    ValueTask QueueBackgroundWorkItemAsync(Func<CancellationToken, Task> workItem);
    ValueTask<Func<CancellationToken, Task>> DequeueAsync(CancellationToken cancellationToken);
}

public class BackgroundTaskQueue : IBackgroundTaskQueue
{
    private readonly Channel<Func<CancellationToken, Task>> _queue;

    public BackgroundTaskQueue(int capacity)
    {
        var options = new BoundedChannelOptions(capacity)
        {
            FullMode = BoundedChannelFullMode.Wait
        };
        _queue = Channel.CreateBounded<Func<CancellationToken, Task>>(options);
    }

    public async ValueTask QueueBackgroundWorkItemAsync(Func<CancellationToken, Task> workItem)
    {
        if (workItem == null) throw new ArgumentNullException(nameof(workItem));
        await _queue.Writer.WriteAsync(workItem);
    }

    public async ValueTask<Func<CancellationToken, Task>> DequeueAsync(CancellationToken cancellationToken)
    {
        return await _queue.Reader.ReadAsync(cancellationToken);
    }
}

第二步:实现后台任务消费者

public class BackgroundTaskWorker : BackgroundService
{
    private readonly ILogger<BackgroundTaskWorker> _logger;
    private readonly IBackgroundTaskQueue _taskQueue;

    public BackgroundTaskWorker(ILogger<BackgroundTaskWorker> logger, IBackgroundTaskQueue taskQueue)
    {
        _logger = logger;
        _taskQueue = taskQueue;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("后台任务服务已启动");

        while (!stoppingToken.IsCancellationRequested)
        {
            try
            {
                var workItem = await _taskQueue.DequeueAsync(stoppingToken);
                await workItem(stoppingToken);
            }
            catch (OperationCanceledException)
            {
                // 忽略任务取消操作
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "处理后台任务时发生异常");
            }
        }

        _logger.LogInformation("后台任务服务已停止");
    }
}

第三步:注册服务到DI容器

在Program.cs中添加:

// 注册任务队列(设置最大队列容量为100)
builder.Services.AddSingleton<IBackgroundTaskQueue>(_ => new BackgroundTaskQueue(100));
// 注册后台任务消费者
builder.Services.AddHostedService<BackgroundTaskWorker>();
// 注册内存缓存用于跟踪任务状态(分布式场景可改用Redis/数据库)
builder.Services.AddMemoryCache();

3. Controller中的业务逻辑处理

这里以批量数据处理为例,判断耗时阈值后分同步/异步处理:

public class DataProcessingController : Controller
{
    private readonly IBackgroundTaskQueue _taskQueue;
    private readonly IMemoryCache _cache;

    public DataProcessingController(IBackgroundTaskQueue taskQueue, IMemoryCache cache)
    {
        _taskQueue = taskQueue;
        _cache = cache;
    }

    [HttpPost]
    public async Task<IActionResult> ProcessLargeData(List<int> dataIds)
    {
        const int threshold = 30000; // 30秒阈值
        // 预估耗时:假设每条数据处理需要10ms
        var estimatedTimeMs = dataIds.Count * 10;

        if (estimatedTimeMs <= threshold)
        {
            // 同步处理短任务
            await ProcessDataSync(dataIds);
            return Ok("数据处理完成!");
        }
        else
        {
            // 生成唯一任务ID
            var taskId = Guid.NewGuid().ToString();
            // 初始化任务状态缓存
            _cache.Set(taskId, new TaskStatus 
            { 
                Id = taskId, 
                Status = "待执行", 
                EstimatedTime = $"{estimatedTimeMs / 1000}秒" 
            }, TimeSpan.FromHours(24));

            // 将任务加入后台队列
            await _taskQueue.QueueBackgroundWorkItemAsync(async token =>
            {
                try
                {
                    // 更新任务状态为处理中
                    _cache.Set(taskId, new TaskStatus 
                    { 
                        Id = taskId, 
                        Status = "处理中", 
                        EstimatedTime = $"{estimatedTimeMs / 1000}秒" 
                    }, TimeSpan.FromHours(24));

                    // 执行实际数据处理逻辑
                    await ProcessDataAsync(dataIds, token);

                    // 更新任务状态为完成
                    _cache.Set(taskId, new TaskStatus 
                    { 
                        Id = taskId, 
                        Status = "已完成", 
                        Result = "数据处理成功",
                        CompletedTime = DateTime.Now
                    }, TimeSpan.FromHours(24));

                    // 发送完成通知(可替换为邮件/站内信)
                    await SendCompletionNotification(taskId);
                }
                catch (Exception ex)
                {
                    // 更新任务状态为失败
                    _cache.Set(taskId, new TaskStatus 
                    { 
                        Id = taskId, 
                        Status = "失败", 
                        ErrorMessage = ex.Message,
                        CompletedTime = DateTime.Now
                    }, TimeSpan.FromHours(24));

                    // 发送失败通知
                    await SendFailureNotification(taskId, ex.Message);
                }
            });

            return Ok(new 
            { 
                Message = $"任务已转入后台执行,预计耗时{estimatedTimeMs/1000}秒,完成后将通知您",
                TaskId = taskId
            });
        }
    }

    // 任务状态查询接口,供前端轮询
    [HttpGet("task-status/{taskId}")]
    public IActionResult GetTaskStatus(string taskId)
    {
        if (_cache.TryGetValue<TaskStatus>(taskId, out var status))
        {
            return Ok(status);
        }
        return NotFound("任务不存在或已过期");
    }

    // 以下为模拟业务方法
    private async Task ProcessDataSync(List<int> dataIds) => await Task.Delay(500);
    private async Task ProcessDataAsync(List<int> dataIds, CancellationToken token) => await Task.Delay(estimatedTimeMs, token);
    private async Task SendCompletionNotification(string taskId) => await Task.CompletedTask;
    private async Task SendFailureNotification(string taskId, string error) => await Task.CompletedTask;
}

// 任务状态模型
public class TaskStatus
{
    public string Id { get; set; }
    public string Status { get; set; }
    public string EstimatedTime { get; set; }
    public string Result { get; set; }
    public string ErrorMessage { get; set; }
    public DateTime? CompletedTime { get; set; }
}

4. 前端交互优化

  • 用户提交数据后,立即展示后台任务提示和任务ID
  • 前端通过轮询/task-status/{taskId}接口获取任务状态,实时更新UI
  • 任务完成/失败时,展示通知弹窗或跳转至结果页面
关键注意事项
  • 数据一致性:如果采用实时监控耗时转后台的方式,要确保已执行的操作要么回滚,要么后台任务具备幂等性(比如通过唯一标识判断数据是否已处理)
  • 任务可靠性:原生BackgroundService在应用重启时会丢失未完成的任务,若需持久化,可改用Hangfire(自带任务持久化、重试机制和可视化界面)
  • 资源控制:限制后台任务的并发数,避免占用过多CPU、内存或数据库连接,影响主应用正常运行
  • 通知可靠性:建议同时支持多种通知方式(邮件+站内信),并允许用户主动查询任务状态

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:08:46