HangFire实战:构建稳定库存同步系统,避开定时任务那些坑
简介.net 8使用HangFire实现库存同步的Demo项目面向电商平台后端开发者演示了京东、天猫、抖音等O2O/B2C场景下库存数据的高效同步方案。项目围绕HangFire后台任务调度、Redis缓存加速与SqlSugar数据库操作完整覆盖了定时同步、队列处理、执行状态监控等关键环节能够帮助开发者快速理解分布式库存服务的核心思路。压缩包共2000个文件以C#源码638个cs和程序集571个dll为主体辅以XML文档、JSON配置文件、CSS/JS前端资源以及NuGet包文件等整体大小139.71MB。目录结构划分清晰包含HangfireServer调度配置、StockServers核心同步逻辑、SkuServers库存单元处理以及独立测试项目便于按模块研读与二次开发。资源目前已有165人学习下载。对于需要借助.net 8原生能力结合HangFire、Redis、SqlSugar构建高可靠库存同步服务的开发者这是一份难得的实战参考既可用于理解技术集成方式也能直接作为项目脚手架进行扩展。 还记得前年给业务线做多平台库存同步的那段时间每天早上打开消息列表之前都要先看一眼定时同步任务有没有报错。库存超卖、断货、客服那边被买家连环问最后所有压力都会落到开发头上。后来我用.NET 8 HangFire重新做了一个库存同步Demo把定时拉取、全量对账、失败重试、任务防重这些基本功理顺整个系统才算真正稳定下来。这篇文章就拿这个Demo当例子聊聊HangFire在库存同步场景里的落地方案和实际踩过的坑适合刚接触HangFire、或者想用后台任务框架替代裸写定时器的朋友参考。1. 为什么是HangFire而不是自己写定时器1.1 库存同步的真实需求很多人以为库存同步就是定时把上游库存数覆盖到本地数据库真做起来完全不是这么回事。我手上这个需求来自一个电商中台场景上游ERP系统维护真实库存下游多个销售渠道需要读取本地库存副本而本地副本不能直接连上游数据库只能通过对方提供的REST API拉数据。这个业务有几个硬性约束。第一API按分页返回每页最多500条所以同步必须处理大量分页请求。第二上游只在商品库存变动后更新UpdatedAt字段没有推送能力所以我们只能靠轮询。第三同步过程不能阻塞主业务接口也不能因为一次网络超时就把后续流程全部卡死。第四凌晨还要对一次全量账防止增量拉取漏掉数据后长期不自知。这些需求叠加起来定时任务不是难在定时这个动作而是难在任务调度、失败重试、执行历史追踪、任务防重这堆基础设施上。自己用BackgroundService写不是不行但调度逻辑、失败恢复、可视化监控全都得自己造轮子风险和工作量都不小。1.2 三种常见实现方案对比我梳理过手头三种可选方案一是裸写BackgroundService二是用Quartz.NET三是用HangFire。这里直接给个对比结论。裸写BackgroundService的优点是依赖少、逻辑完全可控但缺点也最明显没有自带的任务持久化和失败重试机制进程一重启内存里的定时状态就丢了。而且所有执行记录都得自己入库想看某次任务为什么失败得翻日志表体验很差。Quartz.NET是老牌调度库调度能力很强但它在ASP.NET Core里的集成体验比较重。你要自己管IScheduler的单例生命周期、自己写Job的依赖注入解析、自己做执行日志。如果只是做同步任务Quartz会引入不少和业务无关的样板代码。HangFire则是把调度、执行、监控打包在一起。它把任务队列持久化到数据库进程重启后任务不会丢自带Dashboard执行成功失败、耗时多久、重试了几次一眼就能看到失败任务默认自动重试还支持延时任务、延续任务。对于库存同步这种周期性任务 异步任务 失败重试的组合场景HangFire的匹配度非常高。1.3 HangFire解决的核心问题用HangFire解决库存同步本质上是把定时调度、任务队列、执行记录、失败重试这四件事从业务代码里剥离出去。业务代码只关心同步逻辑本身拉数据、比对、写库。至于什么时候触发、触发失败怎么办、执行完怎么查看结果这些都是HangFire的职责。以我的Demo为例里面只用到了HangFire三种任务类型。定时任务RecurringJob负责周期性的增量同步和每天一次的全量对账后台任务BackgroundJob负责手动触发同步请求延续任务Continuation可以在同步完成后发出通知。这三种任务加起来基本覆盖了同步系统的全部流程编排需求。提示HangFire的RecurringJob和Quartz的Cron Trigger在能力上没有本质区别但HangFire把执行日志重试看板都内置了省掉的不只是代码还有后续排查问题的时间。2. Demo设计与接入HangFire2.1 业务假设与核心表为了把Demo讲得具体一点我的业务假设如下本地数据库有一张商品表Products一张库存表InventoryItems还有一张同步记录表SyncLogs。库存表保存上游同步过来的可售库存和锁定库存同步日志则负责记录每一次同步的开始时间、结束时间、状态、处理条数和错误信息。CREATE TABLE InventoryItems ( Id INT IDENTITY PRIMARY KEY, Sku NVARCHAR(50) NOT NULL, Stock INT NOT NULL, LockedStock INT NOT NULL, UpdatedAt DATETIME2 NOT NULL, CONSTRAINT UQ_Sku UNIQUE (Sku) ); CREATE TABLE SyncLogs ( Id INT IDENTITY PRIMARY KEY, SyncType NVARCHAR(20) NOT NULL, StartedAt DATETIME2 NOT NULL, FinishedAt DATETIME2 NULL, Status NVARCHAR(20) NOT NULL, TotalCount INT NOT NULL, FailedCount INT NOT NULL, Message NVARCHAR(MAX) NULL );增量同步还需要记录一个游标就是上一次拉到哪个时间点。我单独建了一张SyncStates表用一行记录增量同步的最后同步时间。上游API通过updatedSince参数过滤变更数据本地记录游标既能减少重复拉取又能保证数据覆盖。这张表结构看起来简单但它是整个增量同步的核心。没有游标每次只能做全量游标记录错了要么漏数据要么每次都拉全量。2.2 模拟上游库存APIDemo本地没有真实的ERP系统所以我用ASP.NET Core Minimal API模拟了一个上游库存接口。这个模拟接口维护了一份内存库存列表支持分页支持根据updatedSince过滤变更数据返回结构固定为{ total, items }。app.MapGet(/api/inventory, (int page, int pageSize, DateTime? updatedSince) { var query _inventory.AsQueryable(); if (updatedSince.HasValue) { query query.Where(x x.UpdatedAt updatedSince.Value); } var total query.Count(); var items query.OrderBy(x x.Sku) .Skip((page - 1) * pageSize) .Take(pageSize) .Select(x new { x.Sku, x.Stock, x.LockedStock, x.UpdatedAt }) .ToList(); return Results.Ok(new { total, items }); });为什么要模拟这一个接口因为很多人在学习同步任务时卡住的不是HangFire本身而是我没有真实上游数据源。用一个内存接口把上游依赖隔离掉整个Demo就能完全本地运行跑通之后再替换成真实HTTP客户端即可。2.3 一个能跑起来的Program.cs配置HangFire接入.NET 8项目需要两个包Hangfire.AspNetCore和Hangfire.SqlServer。如果你的数据库是PostgreSQL可以换成Hangfire.PostgreSql社区的包维护得也还行。Demo里我直接用SQL Server连接字符串指向本地的LocalDB或者开发库。builder.Services.AddHangfire(config { config.SetDataCompatibilityLevel(CompatibilityLevel.Version_180); config.UseSimpleAssemblyNameTypeSerializer(); config.UseRecommendedSerializerSettings(); config.UseSqlServerStorage(builder.Configuration.GetConnectionString(HangfireConnection)); }); builder.Services.AddHangfireServer(options { options.WorkerCount 2; options.SchedulePollingInterval TimeSpan.FromSeconds(5); });这里有个版本细节HangFire 1.8之后UseHangfireServer扩展方法已标记过时推荐直接使用AddHangfireServer。配置里那三个序列化相关的方法不是摆设老版本升级到1.8之后如果不配置启动时会反复收到序列化兼容性警告甚至出现任务反序列化失败的问题后面第5章我会具体讲。在中间件管道里记得加上Dashboardapp.UseHangfireDashboard(/hangfire, new DashboardOptions { DashboardTitle 库存同步任务中心, Authorization new[] { new HangfireDashboardAuthorizationFilter() } });HangfireDashboardAuthorizationFilter需要自己实现IDashboardAuthorizationFilter接口。开发环境可以直接return true生产环境必须校验登录态否则任何能访问到这个路径的人都能手动触发任务、删除任务这不是开玩笑的事。3. 同步任务的主角增量同步、全量对账与手动触发3.1 增量同步游标加分页拉取增量同步是日常跑得最频繁的任务我设置为每5分钟一次。它的逻辑是从SyncStates表读出上一次的同步时间然后循环调用上游分页接口把变更数据一条条Upsert到本地库存表全部完成后更新游标。public class InventoryIncrementalSyncJob { private readonly IDbContextFactoryAppDbContext _contextFactory; private readonly IExternalInventoryApi _api; public InventoryIncrementalSyncJob( IDbContextFactoryAppDbContext contextFactory, IExternalInventoryApi api) { _contextFactory contextFactory; _api api; } [AutomaticRetry(Attempts 3)] public async Task Run(CancellationToken ct) { await using var db await _contextFactory.CreateDbContextAsync(ct); var state await db.SyncStates .FirstOrDefaultAsync(x x.Name incremental, ct); if (state null) { state new SyncState { Name incremental, LastSyncTime DateTime.UtcNow.AddDays(-1) }; db.SyncStates.Add(state); } var page 1; var totalAffected 0; var cursor state.LastSyncTime; const int pageSize 500; while (true) { var batch await _api.GetChangedInventoryAsync(page, pageSize, cursor, ct); if (batch.Items.Count 0) break; foreach (var item in batch.Items) { await UpsertInventoryAsync(db, item, ct); } totalAffected batch.Items.Count; if (page * pageSize batch.Total) break; page; } state.LastSyncTime DateTime.UtcNow; await db.SaveChangesAsync(ct); } }Upsert我用EF Core 7之后提供的ExecuteUpdateAsync实现存在就更新不存在就插入private async Task UpsertInventoryAsync(AppDbContext db, ExternalInventoryDto item, CancellationToken ct) { var exists await db.InventoryItems.AnyAsync(x x.Sku item.Sku, ct); if (exists) { await db.InventoryItems .Where(x x.Sku item.Sku) .ExecuteUpdateAsync(setters setters .SetProperty(x x.Stock, item.Stock) .SetProperty(x x.LockedStock, item.LockedStock) .SetProperty(x x.UpdatedAt, DateTime.UtcNow), ct); } else { db.InventoryItems.Add(new InventoryItem { Sku item.Sku, Stock item.Stock, LockedStock item.LockedStock, UpdatedAt DateTime.UtcNow }); await db.SaveChangesAsync(ct); } }这里有个细节ExecuteUpdateAsync直接生成SQL到数据库执行走的是EF Core的表达式解析复杂场景下不太好用但做Upsert非常顺手。3.2 全量对账任务增量同步虽然高效却不能发现上游删除了某条库存但本地还留着这种问题。所以每天凌晨两点我安排了一个全量对账任务把上游数据完整拉一遍和本地做比较。全量对账的核心不是简单覆盖而是做三件事本地已有但上游没有的SKU做逻辑删除或标记异常本地没有但上游有的SKU做新增两边都有的SKU以最新更新为准。最后还要把比对结果写进SyncLogs出问题时能追溯。public class InventoryFullSyncJob { public async Task Run(CancellationToken ct) { await using var db await _contextFactory.CreateDbContextAsync(ct); var startedAt DateTime.UtcNow; var log new SyncLog { SyncType full, StartedAt startedAt, Status running }; db.SyncLogs.Add(log); try { var page 1; var upstreamSkus new HashSetstring(); const int pageSize 500; while (true) { var batch await _api.GetAllInventoryAsync(page, pageSize, ct); if (batch.Items.Count 0) break; foreach (var item in batch.Items) { upstreamSkus.Add(item.Sku); await UpsertInventoryAsync(db, item, ct); } if (page * pageSize batch.Total) break; page; } await db.InventoryItems .Where(x !upstreamSkus.Contains(x.Sku) !x.IsDeleted) .ExecuteUpdateAsync(setters setters .SetProperty(x x.IsDeleted, true) .SetProperty(x x.UpdatedAt, DateTime.UtcNow), ct); log.Status success; log.FinishedAt DateTime.UtcNow; log.TotalCount upstreamSkus.Count; } catch (Exception ex) { log.Status failed; log.Message ex.Message; throw; } finally { await db.SaveChangesAsync(ct); } } }注意这个任务里有两个不同性质的写库路径循环里的Upsert是逐条执行最后那个ExecuteUpdateAsync是批量执行。全量对账时如果上游有几千个SKU逐条Upsert会比较慢不过Demo阶段能接受。生产环境建议用临时表加MERGE方案后面第6章会提。3.3 手动触发后台Job与Dashboard配合定时任务只是自动化的一种方式实际操作中总会有系统刚对接完想立刻同步一次看看效果的诉求。我在Demo里加了一个手动触发接口通过HangFire的BackgroundJob.Enqueue入队任务接口本身立即返回任务在后台执行。app.MapPost(/api/sync/trigger, (string type) { if (type full) { var jobId BackgroundJob.EnqueueInventoryFullSyncJob( job job.Run(CancellationToken.None)); return Results.Accepted(new { jobId }); } if (type incremental) { var jobId BackgroundJob.EnqueueInventoryIncrementalSyncJob( job job.Run(CancellationToken.None)); return Results.Accepted(new { jobId }); } return Results.BadRequest(未知的同步类型); });为什么不用同步执行接口因为增量同步拉分页时最坏情况可能要几分钟HTTP请求根本等不起如果触发后用户直接关掉页面同步请求就被中断了。而BackgroundJob.Enqueue把任务交给HangFire的Worker线程池接口秒回执行状态在Dashboard里看了一眼就懂。执行过程中如果想给前端加个同步完成通知可以用BackgroundJob.ContinueJobWith在任务成功后追加一个通知任务。这是HangFire延续任务比较典型的使用场景库存同步Demo里虽然不是必须但非常推荐体会一下这个组合。4. 同步任务防翻车幂等、并发锁与失败重试4.1 一个任务跑两次会发生什么很多人写同步任务时容易忽略一个前提任务框架的可重试机制意味着同一个业务逻辑可能会被执行两次甚至更多次。HangFire在任务失败后会自动重试这是好事但如果同步任务本身不幂等重试就会变成灾难。举个具体例子增量同步每5分钟执行一次如果某次执行到一半进程崩溃HangFire会重新入队从头再跑一次。如果这时上次执行已经更新了一部分库存而本次重试没有做Upsert而是直接INSERT就会出现主键冲突如果是全量覆盖逻辑可能把上次刚更新的数据又覆盖回旧值。所以在Demo里所有同步写入用的都是Upsert语义用SKU作为唯一业务键存在就更新不存在就插入。这个设计是幂等的基础。第二个幂等点是游标的更新时机游标必须在所有分页拉取且写库成功之后再更新不能在拉了一半的时候就更新否则漏掉的数据永远不会被补拉。4.2 防止上一轮还没跑完下一轮就启动增量同步设了5分钟一次但一次执行如果遇到上游API慢可能跑10分钟都没结束。这时候下一轮任务又触发了两个任务同时在写同一批SKU虽然Upsert本身不会崩但会造成无意义的重复执行严重时还会把数据库连接池打满。HangFire开源版没有内置的互斥锁。网上常看到的[DisableConcurrentExecution]属性来自HangFire.Pro是收费功能别文档看混了。开源项目要防任务重叠得自己想办法。我在Demo里推荐用数据库分布式锁在SyncStates表旁边加一张SyncLocks表锁记录包含锁名称、持有者、过期时间。任务执行前先尝试获取锁获取成功才继续执行获取失败直接返回等下一轮再来。public interface IDistributedLockService { Taskbool TryAcquireAsync(string key, TimeSpan ttl, CancellationToken ct); Task ReleaseAsync(string key, string owner, CancellationToken ct); }实现上不必做得太重。拿锁就用一条原子SQL当锁不存在或已过期时才插入/更新成功否则认为获取失败。释放锁时校验持有者编号只允许锁的获得者释放。多实例部署时这个方案在数据库层面是安全的。如果你不想手写锁社区也有Hangfire.Mutex之类基于数据库/Redis的互斥实现原理大同小异就是给Job加一个Filter。4.3 失败自动重试与同步日志HangFire默认对失败任务自动重试10次重试间隔递增。对于库存同步这种和外部API打交道的任务10次太多了如果上游接口挂了10次重试会把错误日志刷爆还会造成无意义的流量。我通常会在Job上加[AutomaticRetry(Attempts 3)]把重试次数压到3次。[AutomaticRetry(Attempts 3, OnAttemptsExceeded AttemptsExceededAction.Fail)] public class InventoryIncrementalSyncJob { // ... }需要明白的是AutomaticRetry只对Job抛出的异常生效。如果你在上游返回假数据但HTTP状态码是200时没有主动抛异常HangFire会认为任务成功再多的重试配置也救不了。因此同步任务里一定要做好数据校验返回条数不对、SKU为空、库存为负数都要视为异常。同步日志是另一个关键点。HangFire自己的Dashboard能看任务执行成功失败但业务数据层面还需要记录这次同步影响了哪些SKU、新增了多少、更新了多少、对账发现多少本地缺失。这个日志不是为了给技术看而是为了给业务方对账用。我在SyncLogs里保存了开始时间、结束时间、状态、总条数、失败条数。每次同步结束后更新对应字段出问题时可以直接说凌晨两点的全量对账同步了3200条发现7个本地缺失SKU而不是让业务方自己去猜。5. 在真实环境里最容易踩的四个坑5.1 定时任务时区Cron表达式默认按UTC计算这是所有HangFire新手几乎必踩的坑。RecurringJob里如果直接写字符串Cron表达式比如0 2 * * *不额外指定TimeZoneHangFire默认按UTC计算。你在国内期望的凌晨2点实际会在UTC凌晨2点也就是北京时间上午10点执行。RecurringJob.AddOrUpdateInventoryFullSyncJob( inventory-full-sync, job job.Run(CancellationToken.None), 0 2 * * *, new RecurringJobOptions { TimeZone TimeZoneInfo.Local });注意一个更容易被忽视的细节Cron.Daily(2, 0)这类Cron辅助方法用的是服务器本地时区而字符串Cron表达式用的是RecurringJobOptions.TimeZone。两种写法混用时看起来都是每天2点实际触发时间可能差了好几个小时。我在Demo里统一使用字符串Cron加显式时区避免代码里暗藏两种规则。生产环境如果在Windows服务器上TimeZoneInfo.Local是China Standard Time换成Linux容器后同样代码得到的时区ID可能是Asia/Shanghai。最稳的做法是把时区ID配到配置文件里启动时对Windows和Linux做一次转换而不是写死TimeZoneInfo.Local。5.2 新版HangFire的序列化三件套在HangFire 1.7升级到1.8的过程中序列化行为有变化。官方推荐在配置里显式声明三件套SetDataCompatibilityLevel(CompatibilityLevel.Version_180)、UseSimpleAssemblyNameTypeSerializer()、UseRecommendedSerializerSettings()。这三句话的作用分别是把数据兼容级别切到1.8格式序列化类型时只记录简单程序集名避免程序集版本号变化后反序列化失败使用推荐的JSON序列化设置。如果漏掉其中任何一个可能出现两类问题——启动后反复出现序列化兼容性警告或者之前入队的旧任务在反序列化时直接报错。特别是用IRecurringJobManager注册任务的生产环境任务定义已经持久化在数据库里。程序集升级后如果类型序列化格式不兼容老任务可能无法反序列化。配置三件套能减免大部分这类问题。5.3 Job里面的DbContext生命周期HangFire的Worker是并发执行的默认WorkerCount会按机器和配置设置通常不小。这意味着多个Job可能同时在跑如果Job里注入的是Scoped的DbContext而Job类型的生命周期是Transient或Singleton很容易出现上下文被多个任务共用、DbContext线程安全问题。我建议在Job里使用IDbContextFactoryAppDbContext每次需要数据库操作时再创建独立的上下文await using var db await _contextFactory.CreateDbContextAsync(ct);这样每个任务、甚至任务里的每个批次都拥有独立的上下文彼此不干扰用完即释放。这个写法在普通API请求里无所谓在后台任务这种长生命周期场景里是必须养成的习惯。还有一点长时间运行的任务不要在一个DbContext里做几万次查询和更新上下文跟踪的实体越来越多内存和性能都会出问题。分批操作时尽量让一个上下文只负责一个批次的读写。5.4 Dashboard不要裸奔HangFire的Dashboard功能强大但也非常危险。如果你把服务部署到公网Dashboard又没有鉴权那么任何访问者都能看到你的任务定义、执行记录、失败明细还能手动触发任务。这等于把同步系统的控制面板直接暴露给所有人。Demo里我写了一个极简的授权过滤器public class HangfireDashboardAuthorizationFilter : IDashboardAuthorizationFilter { public bool Authorize(DashboardContext context) { var httpContext context.GetHttpContext(); return httpContext.User.Identity?.IsAuthenticated true; } }开发环境你可以先返回true一旦上了生产务必接上你现有的登录态校验或者至少限定内网IP访问。Dashboard的路径也不要使用默认的/hangfire改一个不容易被猜到的路径能挡掉一部分扫描流量。6. Demo到生产还差哪些6.1 大库存量下的写入策略Demo里用EF Core逐条Upsert几百上千个SKU完全没问题但库存量到了十万级别逐条ExecuteUpdateAsync会变得非常慢。生产环境我见得比较多的是两种方案第一种是把增量数据先写进临时表然后用一条MERGE语句批量合并到正式表第二种是使用SqlBulkCopy直接把数据灌入临时表再做合并。这两种方案的速度比逐条更新高一个数量级。使用SqlBulkCopy时需要注意目标表结构要和DataTable列对应好批量操作前关闭表上的触发器或索引维护窗口否则大量索引重建会拖慢写入。6.2 多实例部署的分布式互斥Demo里的数据库锁可以防止单实例下任务重叠但生产环境往往是多实例部署HangFire会同时启动多个Server。RecurringJob本身由HangFire内部协调不会重复入队但我们的业务锁必须能在多个实例间互斥。这时数据库锁一样可用但更常见的方案是换成Redis分布式锁比如StackExchange.Redis配合RedLock或者直接用Medallion.Threading这类封装好的库。锁的粒度要精确到任务名比如inventory:full:lock和inventory:incremental:lock分开避免全量同步阻塞增量同步这种不必要的串行。6.3 同步质量的可观测性有了HangFire Dashboard之后单个任务的执行情况是清晰了但整个链路的健康度还需要额外关注。我在生产环境会加两个维度的监控一个是HangFire任务维度的监控失败率超过阈值就告警另一个是业务数据维度的监控定时跑一个对账任务比对上游和本地库存上一天的差异条数差异超过阈值就推告警。第二个维度容易被忽略但价值很高。因为任务框架只能告诉你任务执行成功不能告诉你同步结果对不对。只有业务层面的校验才能发现增量游标漏数据、上游字段映射错了这类隐蔽问题。如果你不想为了告警引入一套复杂系统也可以把HangFire的失败任务信息转发到钉钉或企业微信机器人再配合SyncLogs表的数据质量统计基本够用。常用做法是写一个JobFailedHandler在HangFire的任务状态变为Failed时触发告警推送成本低见效快。我现在的习惯是不管项目大小先把Dashboard、失败重试、同步日志这三件套配齐再谈业务功能。等哪一天任务真的失败了你会发现HangFire已经帮你在界面里把失败现场还原了出来这种体验比自己翻日志舒服太多。本文还有配套的精品资源点击获取