.NET 后台服务实战:BackgroundService + PeriodicTimer + IAsyncEnumerable 实现定时任务同步

在 ASP.NET Core 项目中,除了处理 HTTP 请求的接口逻辑,我们经常需要一些”后台常驻”的能力:定时同步配置、心跳检测、队列消费、数据清理等。.NET 提供了 IHostedService / BackgroundService 这一套托管服务机制,让我们可以用几乎与业务代码一致的写法来编写长期运行的后台任务。

本文基于 SwitchData 项目中的 JobWorkerService 实现,讲解如何组合使用 BackgroundServicePeriodicTimerIAsyncEnumerable,构建一个可取消、可释放、资源友好的后台定时任务同步服务。

一、为什么需要 BackgroundService

在没有托管服务之前,常见的后台任务写法有两种:

  • 在启动代码里直接 Task.Run 一个无限循环;
  • 在控制器或服务里埋一个”定时触发”的逻辑。

这两种方式都能跑,但存在几个问题:

  1. 生命周期不可控:应用关闭时,任务可能还在执行,导致数据不一致或资源泄漏;
  2. 取消令牌缺失:没有统一的 CancellationToken,任务无法感知关闭信号;
  3. 依赖注入不方便:静态类或手动 new 出来的任务难以使用注册在 DI 容器中的服务;
  4. 异常处理随意:缺少统一的日志和重试策略。

BackgroundService 是 .NET 对 IHostedService 的封装,它解决了上述大部分问题。我们只需要重写 ExecuteAsync(CancellationToken stoppingToken),框架就会负责:

  • 在应用启动后调度执行;
  • 传入取消令牌;
  • 在应用关闭前等待任务完成或超时;
  • 与 DI 容器无缝集成。
public class JobWorkerService : BackgroundService
{
    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        // 后台任务主逻辑
    }
}

二、PeriodicTimer:比 Task.Delay 更优雅的定时器

在传统写法中,“每 30 秒执行一次”通常这样实现:

while (true)
{
    DoWork();
    await Task.Delay(TimeSpan.FromSeconds(30), cancellationToken);
}

这种方式简单直接,但有一个小问题:DoWork() 本身的执行时间会被累加到周期里。如果任务执行了 10 秒,那么下一次实际触发间隔就是 40 秒而不是 30 秒。对于心跳检测这类”必须固定间隔触发”的场景,这并不理想。

.NET 6 引入了 PeriodicTimer,它专门用于固定周期的等待:

using var timer = new PeriodicTimer(TimeSpan.FromSeconds(30));
while (await timer.WaitForNextTickAsync(stoppingToken))
{
    SendHeartBeat();
}

WaitForNextTickAsync 会在每个周期点返回 true。如果当前这一次的周期点已经错过,它会立即返回,不会无限等待。更重要的是,它会自动累积周期,不会因为单次任务耗时而拉长整体间隔。

在 SwitchData 项目中,心跳检测就选用了 PeriodicTimer

_heartBeatTimer = new PeriodicTimer(TimeSpan.FromSeconds(30));
while (await _heartBeatTimer.WaitForNextTickAsync(stoppingToken))
{
    SendHeartBeat();
}

三、IAsyncEnumerable:流式读取数据库记录

后台任务经常需要批量处理数据。最 intuitive 的做法是一次性把所有记录查出来:

var list = await db.QueryAsync<TaskEntity>("SELECT * FROM NeTask");
foreach (var item in list)
{
    Process(item);
}

如果记录数只有几百条,这没问题。但当数据量达到万级、十万级时,一次性加载会带来两个风险:

  1. 内存占用高:所有记录都要驻留在内存;
  2. 延迟大:必须等全部数据返回后,处理逻辑才能开始。

IAsyncEnumerable<T> 提供了一种”流式”的异步枚举方式。配合 await foreach,可以在数据从数据库逐行读取的同时进行处理,实现真正的”边读边处理”。

SwitchData 的数据访问层大量使用了这种模式:

public static async IAsyncEnumerable<NeTaskView> GetListAsync(
    bool? enabled = null,
    ExecutionStatus? status = null,
    [EnumeratorCancellation] CancellationToken cancellationToken = default)
{
    var metadata = MetadataCache.GetTableMetadata<NeTaskView>();
    string selectClause = $"SELECT * FROM [{metadata.DbTableName}]";

    using var command = BuildCommand(selectClause, enabled, status);
    await using var reader = await database.ExecuteReaderAsync(command, cancellationToken);

    while (await reader.ReadAsync(cancellationToken))
    {
        yield return ObjectMapper.MapFrom<NeTaskView>(reader);
    }
}

[EnumeratorCancellation] 是关键:它把调用方传入的 CancellationToken 绑定到异步枚举器上,这样 await foreach 内部调用 MoveNextAsync 时,也能响应取消请求。

JobWorkerService 中,我们就是利用这个流式接口来同步所有定时任务:

await foreach (var task in NeTaskService.GetListAsync(cancellationToken: cancellationToken))
{
    totalCount++;
    if (task.Enabled)
    {
        enabledCount++;
        _scheduler.EnableSchedule(task.Id, task.JobId, task.CronExpression);
    }
    else
    {
        _scheduler.DisableSchedule(task.JobId);
    }
}

四、JobWorkerService 完整实现解析

把上面三个技术点组合起来,就得到了 SwitchData 中的 JobWorkerService

public class JobWorkerService : BackgroundService
{
    private readonly NeTaskHangfireScheduler _scheduler;
    private readonly string _workerId = $"{Environment.MachineName}_{Guid.NewGuid():N}";
    private PeriodicTimer _heartBeatTimer;

    public JobWorkerService(NeTaskHangfireScheduler scheduler)
    {
        _scheduler = scheduler;
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        Logger.Info($"定时任务Worker启动,WorkerId:{_workerId}");

        // 启动心跳检测
        _heartBeatTimer = new PeriodicTimer(TimeSpan.FromSeconds(30));

        // 同步定时任务:作为一次性后台任务启动
        _ = Task.Run(() => SyncJobsAsync(stoppingToken), stoppingToken);

        // 主循环:周期性心跳
        while (await _heartBeatTimer.WaitForNextTickAsync(stoppingToken))
        {
            SendHeartBeat();
        }

        Logger.Info("定时任务Worker停止");
    }

    private async Task SyncJobsAsync(CancellationToken cancellationToken)
    {
        try
        {
            int totalCount = 0;
            int enabledCount = 0;

            await foreach (var task in NeTaskService.GetListAsync(cancellationToken: cancellationToken))
            {
                totalCount++;
                if (task.Enabled)
                {
                    enabledCount++;
                    _scheduler.EnableSchedule(task.Id, task.JobId, task.CronExpression);
                }
                else
                {
                    _scheduler.DisableSchedule(task.JobId);
                }
            }

            Logger.Info($"同步定时任务成功:启用 {enabledCount} 个,停止 {totalCount - enabledCount} 个");
        }
        catch (Exception ex)
        {
            Logger.Error("同步定时任务失败", ex);
        }
    }

    private void SendHeartBeat()
    {
        // 心跳上报逻辑
    }

    public override async Task StopAsync(CancellationToken cancellationToken)
    {
        _heartBeatTimer?.Dispose();
        await base.StopAsync(cancellationToken);
    }
}

4.1 WorkerId 的设计

_workerId 由机器名和 GUID 组成,用于在多实例部署时区分不同的 Worker。这对于日志排查、分布式锁设计、任务去重都很有帮助。

4.2 多任务并发的设计取舍

ExecuteAsync 中同时做了两件事:

  1. 启动 SyncJobsAsync 作为一次性任务;
  2. 进入心跳循环。

这里用 _ = Task.Run(...) 把同步任务”fire and forget”地放到后台。注意,这并不是最佳实践——如果 SyncJobsAsync 抛出未捕获的异常,这个异常会被吞掉(虽然方法内部已经 try-catch)。更稳妥的做法是把同步任务也纳入生命周期管理,或者使用 Task.WhenAll 等待两个任务:

protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
    _heartBeatTimer = new PeriodicTimer(TimeSpan.FromSeconds(30));

    var syncTask = SyncJobsAsync(stoppingToken);
    var heartbeatTask = HeartbeatLoopAsync(stoppingToken);

    await Task.WhenAll(syncTask, heartbeatTask);
}

private async Task HeartbeatLoopAsync(CancellationToken stoppingToken)
{
    while (await _heartBeatTimer.WaitForNextTickAsync(stoppingToken))
    {
        SendHeartBeat();
    }
}

4.3 取消与资源释放

PeriodicTimer 实现了 IDisposable,在 StopAsync 中需要主动释放:

public override async Task StopAsync(CancellationToken cancellationToken)
{
    _heartBeatTimer?.Dispose();
    await base.StopAsync(cancellationToken);
}

另外,ExecuteAsync 里的 while 循环依赖 stoppingToken。当应用关闭时,WaitForNextTickAsync 会抛出 OperationCanceledException,循环自然退出。

五、服务注册与项目结构

后台服务写好后,需要在 Program.cs 中注册为托管服务:

builder.Services.AddHostedService<JobWorkerService>();
builder.Services.AddSingleton<NeTaskHangfireScheduler>();

SwitchData 把后台任务独立成一个 SwitchData.Job 项目,专门负责 Hangfire 的调度与 Worker 心跳。这种拆分有几个好处:

  • 职责清晰:Web API 项目只处理 HTTP 请求,Job 项目只处理后台任务;
  • 独立扩展:在高并发场景下,可以单独扩容 Job 实例;
  • 故障隔离:后台任务的异常不会影响 API 接口的可用性。

整个同步流程可以用下图表示:

flowchart LR
    A[JobWorkerService 启动] --> B{启动心跳循环}
    A --> C[SyncJobsAsync 同步任务]
    C --> D[GetListAsync IAsyncEnumerable]
    D --> E[逐条读取 NeTask]
    E --> F{Enabled?}
    F -->|是| G[EnableSchedule 启用 Hangfire]
    F -->|否| H[DisableSchedule 停用 Hangfire]
    B --> I[每 30s SendHeartBeat]

六、常见问题与最佳实践

6.1 不要阻塞 ExecuteAsync

ExecuteAsync 应该尽快进入异步等待状态。如果里面有同步阻塞调用(如 Thread.Sleep),会占用线程池线程,影响托管服务的正常调度。

6.2 处理 OperationCanceledException

当应用关闭时,stoppingToken 会被触发,此时 WaitForNextTickAsyncReadAsync 等操作都会抛出 OperationCanceledException。这是正常行为,不需要记录为错误:

catch (OperationCanceledException)
{
    // 正常关闭,忽略
}
catch (Exception ex)
{
    Logger.Error("后台任务异常", ex);
}

6.3 避免在 ExecuteAsync 中直接 await 长时间同步代码

如果必须执行 CPU 密集型操作,可以把它放到 Task.Run 中,并传入 stoppingToken

await Task.Run(() => HeavyComputation(stoppingToken), stoppingToken);

6.4 数据库连接的及时释放

在使用 IAsyncEnumerable 时,using 语句会绑定到枚举器的生命周期。调用方应该使用 await foreach 而不是先 ToList() 再遍历,否则就失去了流式读取的意义。

七、总结

BackgroundService + PeriodicTimer + IAsyncEnumerable 是 .NET 6+ 中一套非常顺手的高性能后台任务组合:

  • BackgroundService 提供了标准的托管服务生命周期;
  • PeriodicTimer 让固定周期的心跳/轮询更精确、更省资源;
  • IAsyncEnumerable 让大数据量的数据库同步不再一次性占用大量内存。

在 SwitchData 项目中,这套组合被用来实现定时任务 Worker 的启动同步与心跳检测。理解它们的协作方式,可以帮助我们写出更稳定、更可维护的后台服务。