ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

如何用 Go 后台任务队列 River 把耗时逻辑全部交给异步

如何用 Go 后台任务队列 River 把耗时逻辑全部交给异步 如何用 Go 后台任务队列 River 把耗时逻辑全部交给异步【免费下载链接】riverFast and reliable background jobs in Go项目地址: https://gitcode.com/gh_mirrors/river/river大促当晚订单量是平日的十倍。如果你在下单接口里同步发邮件、同步生成对账单接口平均耗时会从 80ms 一路涨到好几秒前端页面开始转圈用户疯狂点重试——这个画面做过 Web 后端的人都不陌生。把耗时操作挪出请求链路、放到后台慢慢执行是所有这类系统的刚需。而 River 正是为 Go 语言打造的高性能后台作业处理框架任务先落进数据库再由独立的工作节点异步消化接口只负责接单不再承担脏活累活。先让概念有画面感作业、队列、工作节点理解 River 之前不妨把整个系统想象成一家外卖店**作业Job**就是一张订单小票。在代码里它只是一个实现了Kind()方法的普通 struct用来描述要做什么、参数是什么。**队列Queue**相当于店里的分单台。不同队列可以分给不同窗口处理例如emails队列专门处理邮件、report队列专门跑报表互不干扰。**工作节点Worker**是后厨里真正干活的店员。它持续从队列里取单、执行干完一单再取下一单。这套画面记在心里后面所有机制都能对号入座。一次任务的完整旅程从入队到收尾那么问题来了一张订单从提交到执行完River 内部究竟让它经历了什么我们可以用提交 → 排队 → 取单 → 执行 → 收尾五步来串第一步提交。调用客户端的Insert方法把作业交出去参数会被序列化JSON后写入数据库并立刻返回这条作业的 ID 和当前状态。整个过程只有一次数据库写入请求链路轻得几乎感觉不到。第二步排队。作业在库里并不会立刻被取走。如果设了延迟执行时间它会先处于scheduled状态后台的调度服务实现于 internal/maintenance/job_scheduler.go会定时扫描把到点的作业翻成available就绪状态周期性作业的重复投递也由它负责。第三步取单。工作节点持续从库里捞处于就绪状态的作业。为了避免多个节点同时抢到同一单River 在取单和状态推进时用了乐观锁——通过对比版本号或尝试次数来确保同一时刻只有一个执行者减少锁等待。第四步执行。取到单后River 会起一个 goroutine 运行你写的Work逻辑。跑成功作业标记为completed跑失败则按重试策略安排进入retryable状态等待下一次尝试。第五步收尾。完成或废弃的作业不会永远躺在表里。维护服务里的清理器internal/maintenance/job_cleaner.go会按保留期定期删除过期记录控制表体积保证查询效率。五分钟跑通最简闭环一个排序任务的完整代码把流程讲得再清楚也不如一段能直接跑起来的代码。下面用 PostgreSQL 驱动riverpgxv5演示最简闭环定义参数 → 注册 Worker → 创建客户端 → 提交作业。// 1. 定义作业参数描述要做什么Kind 用于标识作业类型 type SortArgs struct { Strings []string json:strings } func (SortArgs) Kind() string { return sort } // 2. 定义工作节点真正干活的地方 type SortWorker struct { river.WorkerDefaults[SortArgs] // 嵌入默认实现省去写一堆模板方法 } func (w *SortWorker) Work(ctx context.Context, job *river.Job[SortArgs]) error { sort.Strings(job.Args.Strings) fmt.Println(sorted:, job.Args.Strings) return nil // 返回 nil 表示成功返回 error 则触发重试 } // 3. 组装客户端Driver 负责对接数据库Config 负责描述并发与队列 func main() { ctx : context.Background() dbPool, _ : pgxpool.New(ctx, postgres://...) defer dbPool.Close() workers : river.NewWorkers() river.AddWorker(workers, SortWorker{}) riverClient, err : river.NewClient(riverpgxv5.New(dbPool), river.Config{ Queues: map[string]river.QueueConfig{ river.QueueDefault: {MaxWorkers: 100}, // 该队列最多 100 个并发 }, Workers: workers, }) if err ! nil { panic(err) } // 4. 启动客户端worker 开始从库里取单 if err : riverClient.Start(ctx); err ! nil { panic(err) } // 5. 提交作业参数会被序列化后写入数据库 if _, err : riverClient.Insert(ctx, SortArgs{ Strings: []string{tiger, whale, bear}, }, nil); err ! nil { panic(err) } // 到这里接口就可以返回了排序会在后台完成 }创建客户端的入口是NewClient见 client.go它需要两个东西驱动Driver和配置Config。驱动决定了作业存在哪个数据库里目前官方提供 PostgreSQLriverpgxv5、SQLiteriversqlite以及面向通用 SQL 的riverdatabasesql实现接口保持一致切换后端基本不用改业务代码。遇到问题River 是怎么接招的入门之后真实业务里会有各种刁钻场景River 大多已有现成解法问题一调第三方接口偶尔超时任务失败了怎么办River 的默认重试策略采用指数退避失败后的等待间隔大致按尝试次数^4增长第一次失败 1 秒、第二次约 16 秒、第三次约 1 分 21 秒直到达到最大尝试次数仍失败才把作业标记为discarded废弃。重试相关逻辑在 retry_policy.go 中定义你可以自定义策略也可以按单条作业通过InsertOpts.MaxAttempts调整上限。问题二想延迟执行比如下单后 30 分钟再提醒支付提交时给InsertOpts.ScheduledAt传入未来时间即可作业会先进入 scheduled 状态到点由调度器自动放行保证不会提前执行。问题三队列里任务太多怎么控制资源Config.Queues中每个队列的MaxWorkers就是它的并发上限River 会按这个数字起对应数量的 goroutine 并行取单执行充分利用多核 CPU同时每条作业还支持 1–4 的优先级1 最高保证重要任务先被取走。问题四进程重启会不会丢任务会不会多个实例抢活作业持久化在数据库里重启后继续执行天然不丢多实例部署时River 通过领导者选举internal/leadership/elector.go保证同一时刻只有一个实例在跑调度与维护类任务避免重复操作和单点故障。新手常踩的三个坑坑一客户端没启动就插入作业。作业确实入库了但 worker 没开始取单任务永远卡在队列里。记住顺序先Start再Insert。坑二重试配置走极端。MaxAttempts设成 1瞬时网络抖动就把任务打成了 discarded数据后续还得人工补设得太大故障期间任务会反复重试打爆日志。建议按业务容忍度设置并配合 ErrorHandler 做好失败记录。坑三Work 方法忽略 context 取消。优雅停机时 River 会通过取消 context 通知正在执行的任务如果你的Work里是死循环或阻塞调用、完全不看ctx.Done()进程可能迟迟退不干净。写网络请求、定时器时记得用select监听 ctx。收个尾River 适合所有不想自研异步任务系统的 Go 项目接口要快、任务不能丢、失败要自动重试。它把数据库既当存储又当调度中枢思路简单可靠性和性能却都不含糊。想继续深入建议从 docs/development.md 看起再顺着维护服务目录把调度器、清理器逐个翻一遍你对后台作业系统的理解会上一个新台阶。【免费下载链接】riverFast and reliable background jobs in Go项目地址: https://gitcode.com/gh_mirrors/river/river创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表