Loco Workers 后台任务完整指南:四种队列后端、Worker 模式与任务生命周期管理
2026/9/16 15:54:43 网站建设 项目流程

Loco Workers 后台任务完整指南:四种队列后端、Worker 模式与任务生命周期管理

【免费下载链接】loco🚂 🦀 The one-person framework for Rust for side-projects and startups项目地址: https://gitcode.com/GitHub_Trending/lo/loco

本文是 Loco(Rust 单人全栈框架)后台任务(Workers)模块的完整实战指南。围绕 docs-site/content/docs/processing/workers.md 展开,你将掌握 Loco 的四种后台任务运行方式(Redis / Postgres / SQLite 队列与 Tokio 同进程异步)、三种 Worker 模式(BackgroundQueue/ForegroundBlocking/BackgroundAsync)的配置与取舍、perform_later强类型入队、Worker 标签过滤、cargo loco generate worker脚手架,以及通过 CLI 对任务队列进行 cancel / tidy / purge / dump / import / requeue 的全生命周期管理。

Loco 后台任务的抽象设计:一套 API,四种后端

Loco 为后台任务提供了统一的抽象,你的业务代码只与BackgroundWorkertrait 和perform_later打交道,而无需关心底层队列实现——这与 Rails 的ActiveJob思路一致:通过一次配置变更即可切换队列后端,代码零改动(见 src/bgworker/mod.rs)。

后台任务有四种可用后端:

  • Redis backed:由独立的 Redis 服务承载任务队列,适合多进程、跨机器分发负载;
  • Postgres backed:把任务表放进 Postgres 数据库,与业务数据同库管理;
  • SQLite backed:任务存储在 SQLite 中,适合本地开发与单机部署;
  • Tokio-async 同进程模式:不依赖任何外部存储,任务在当前服务器进程的 Tokio 异步池中直接执行。

对应的三个队列特性开关在根目录 Cargo.toml 中以 feature 形式提供:bg_redisbg_pgbg_sqlt(分别依赖redissqlxcrate)。在源码中,Queue枚举(见 src/bgworker/mod.rs)统一封装了三种队列提供方与None(无队列)四种状态,enqueueregisterrunsetuppingclear等方法对三种后端做了一致的分发处理。

三种 Worker 模式:BackgroundQueue / ForegroundBlocking / BackgroundAsync

config/<environment>.yaml中通过workers.mode指定 Worker 的运行模式。三种模式的定义见 src/config.rs:

模式说明是否依赖外部队列
BackgroundQueue任务进入 Redis/Postgres/SQLite 队列,由独立 Worker 进程异步消费,默认模式需要
ForegroundBlocking任务在前台同步执行、阻塞直到完成不需要
BackgroundAsync任务在同一个进程内以 Tokio 异步任务执行不需要

配置示例(写入config/<environment>.yaml):

# Worker Configuration workers: # specifies the worker mode. Options: # - BackgroundQueue - Workers operate asynchronously in the background, processing queued. # - ForegroundBlocking - Workers operate in the foreground and block until tasks are completed. # - BackgroundAsync - Workers operate asynchronously in the background, processing tasks with async capabilities. mode: BackgroundQueue

三点关键差异需要理解:

  1. BackgroundQueue是默认模式(源码中WorkerMode标注了#[default]),它要求你在配置中提供queue段,否则perform_later会打印错误日志“background queue is selected, but queue was not populated in context”;
  2. ForegroundBlockingBackgroundAsync都不需要 Redis/Postgres/SQLite,直接在进程内执行;
  3. ForegroundBlocking阻塞调用方直到任务完成BackgroundAsync则通过tokio::spawn把任务抛到异步池后立即返回。从perform_later的实现(src/bgworker/mod.rs)可以看到三种模式的分流逻辑:BackgroundQueue走队列入队,ForegroundBlocking直接同步performBackgroundAsynctokio::spawn异步执行。

BackgroundAsyncBackgroundQueue的核心区别在于:前者把任务存进 Redis/Postgres/SQLite 并交给独立进程消费,后者则在当前进程的内存异步池中运行

配置队列后端:Redis / Postgres / SQLite 参数全解

workers.mode设为BackgroundQueue时,需要在配置中声明queue段。三种后端的配置结构定义在 src/config.rs,下面给出完整参数与默认值说明。

Redis 队列后端

queue: kind: Redis # Redis connection URI. uri: "{{ get_env(name="REDIS_URL", default="redis://127.0.0.1") }}" # Dangerously flush all data. dangerously_flush: false # represents the number of tasks a worker can handle simultaneously. num_workers: 2

对应RedisQueueConfig(src/config.rs),除上表字段外还支持:

  • queues:可选,声明自定义队列名列表,第一个队列优先级最高,可用于建模优先级队列(Option<Vec<String>>,默认None);
  • num_workers:单个 Worker 进程可同时处理的任务数,默认值 2(由num_workers()函数提供,见 src/config.rs)。

Postgres 队列后端

queue: kind: Postgres # Postgres Queue connection URI. uri: "{{ get_env(name="PGQ_URL", default="postgres://localhost:5432/mydb") }}" # Dangerously flush all data. dangerously_flush: false # represents the number of tasks a worker can handle simultaneously. num_workers: 2

对应PostgresQueueConfig(src/config.rs),完整字段与默认值如下:

字段说明默认值
uriPostgres 连接串必填
dangerously_flushtrue时启动即清空全部任务false
enable_logging是否启用 SQL 日志false
max_connections连接池最大连接数20
min_connections连接池最小连接数1
connect_timeout连接超时(毫秒)500
idle_timeout空闲超时(毫秒)500
poll_interval_sec轮询新任务的间隔(秒)1
num_workers每个进程并发处理任务数2

SQLite 队列后端

queue: kind: Sqlite # SQLite Queue connection URI. uri: "{{ get_env(name="SQLTQ_URL", default="sqlite://loco_development.sqlite?mode=rwc") }}" # Dangerously flush all data. dangerously_flush: false # represents the number of tasks a worker can handle simultaneously. num_workers: 2

对应SqliteQueueConfig(src/config.rs),字段与 Postgres 版一致(max_connections/min_connections/connect_timeout/idle_timeout/poll_interval_sec/num_workers),默认值相同,其中poll_interval_sec由独立的sqlt_poll_interval()提供、默认 1 秒。

关于dangerously_flushconverge函数(src/bgworker/mod.rs)在启动初始化阶段检查该标志,为true时会对所选后端执行queue.clear(),把队列中所有任务清空,因此命名为“危险”选项,生产环境请谨慎开启。

选择 Async 还是 Queue:分发负载与跨服务器扩展

创建新应用时默认会选择async配置,即任务在 Tokio 异步池中运行——这适合单进程、无需额外基础设施的场景:任务与 Web 服务共用一个进程,获得真正的“后台处理”。

当你需要把负载分发到多台服务器时,则应切换到独立的队列进程:

  1. workers.mode改为BackgroundQueue(见上文配置);
  2. 配置 Redis / Postgres / SQLite 三选一的queue段;
  3. 单独启动一个或多个 Worker 进程来消费任务(见下一节)。

队列提供方的创建在create_queue_provider(src/bgworker/mod.rs)中完成:仅当workers.mode == BackgroundQueue且存在queue配置时才会构建对应的 provider 并注入AppContextctx.queue_provider);如果模式不是BackgroundQueue或没有配置队列,则返回None

运行 Worker 进程:start 命令的三种形态

cargo loco start的子命令决定了后台任务在哪里执行,完整参数如下(源码见 src/cli.rs):

Usage: demo_app start [OPTIONS] Options: -w, --worker [<WORKER>...] Start worker. Optionally provide tags to run specific jobs (e.g. --worker=tag1,tag2) -s, --server-and-worker start same-process server and worker -a, --all start server, worker, and scheduler in the same process

场景一:独立 Worker 进程。当你配置了真正的 Redis/Postgres/SQLite 队列,希望单独开一个进程专职消费后台任务时使用--worker。此时每台服务器可以跑一个 Worker 进程,主 Web/API 服务用cargo loco start正常启动:

$ cargo loco start --worker # starts a standalone worker job executing process $ cargo loco start # starts a standalone API service or Web server, no workers

场景二:同进程 Server + Worker。当你使用async后台 Worker 时,任务会随当前服务器进程一起执行,应使用-s

$ cargo loco start --server-and-worker # both API service and workers will execute

场景三:全量启动。还可使用--all让服务器、Worker 与调度器(scheduler)在同一进程内运行。三个选项互斥(conflicts_with_all),CLI 层面对其做了校验;--worker传参时以逗号分隔、支持 0 个或任意多个标签(value_delimiter = ',')。此外cargo loco watch也支持-w/--worker-s/--server-and-worker,方便开发时热重载。

Worker 标签过滤:为特定任务建立专职消费者

Loco 支持基于标签的任务过滤:你可以创建只处理特定类型任务的专职 Worker,非常适合负载分发或为资源密集型任务(如报表生成、数据分析)单独开进程。

启动 Worker 时通过--worker指定要处理的标签:

# Start a worker that only processes jobs with no tags $ cargo loco start --worker # Start a worker that only processes jobs with the "email" tag $ cargo loco start --worker email # Start a worker that processes jobs with either "report" or "analytics" tags $ cargo loco start --worker report,analytics

关于标签匹配的重要规则(原文档明确约定):

  1. 不带标签启动的 Worker(cargo loco start --worker只处理没有标签的任务
  2. 带标签启动的 Worker 只处理至少匹配其中一个标签的任务;
  3. --all--server-and-worker模式不支持标签过滤,只会处理无标签任务;
  4. 标签区分大小写

标签从 CLI 传入后,最终会透传到Queue::run(tags)(src/bgworker/mod.rs),进而下发给对应后端的 JobRegistry 进行按标签消费。

在代码中创建后台任务:强类型参数的 perform_later

使用 Worker 的核心动作是把任务加入队列。在控制器等任何持有AppContext的地方,调用 Worker 的perform_later即可:

// .. in your controller .. DownloadWorker::perform_later( &ctx, DownloadWorkerArgs { user_guid: "foo".to_string(), }, ) .await

与 Rails/Ruby 不同,Rust 让你享受强类型任务参数DownloadWorkerArgs会被序列化(Redis 后端序列化为可存储格式,Postgres/SQLite 后端经serde_json::to_value序列化)后推入队列,消费端再用同一结构反序列化,编译期即可发现参数不匹配问题。

perform_laterBackgroundWorkertrait 的静态方法(src/bgworker/mod.rs),其内部按当前workers.mode自动分流:

  • BackgroundQueue:把任务以class_name()作为标识、连同队列与标签一起enqueuectx.queue_provider
  • ForegroundBlocking:立即Self::build(ctx).perform(args)同步执行;
  • BackgroundAsynctokio::spawn异步执行,出错时记录worker failed to perform job日志。

给任务打标签:定义 Worker 的 tags()

入队时可以为任务附带标签,之后只有匹配标签的 Worker 才会处理它。标签在 Worker 的tags()方法中声明:

// To create a job with a tag, define the tags in your worker: struct DownloadWorker; #[async_trait] impl BackgroundWorker<DownloadWorkerArgs> for DownloadWorker { // Define tags for this worker fn tags() -> Vec<String> { vec!["download".to_string(), "network".to_string()] } // ... other implementation details } // When you call perform_later, the job will automatically be tagged DownloadWorker::perform_later(&ctx, args).await?;

源码层面(src/bgworker/mod.rs),tags()默认返回空向量;perform_later在入队时会把非空的标签列表转为Some(tags)传给enqueue

创建新 Worker:实现 trait 并注册到全局任务处理器

新增 Worker 意味着两件事:编写接受参数并执行任务的业务逻辑,以及让 Loco 认识它并注册进全局任务处理器。

第一步:在workers/目录添加 Worker。DownloadWorker为例(完整模板可参考 loco-new/base_template/src/workers/downloader.rs):

#[async_trait] impl BackgroundWorker<DownloadWorkerArgs> for DownloadWorker { fn build(ctx: &AppContext) -> Self { Self { ctx: ctx.clone() } } // Optional: Define tags for this worker fn tags() -> Vec<String> { vec!["download".to_string()] } async fn perform(&self, args: DownloadWorkerArgs) -> Result<()> { println!("================================================"); println!("Sending payment report to user {}", args.user_guid); // TODO: Some actual work goes here... println!("================================================"); Ok(()) } }

第二步:在app.rsconnect_workers中注册。新项目模板的app.rs.t中已经预留了钩子(见 loco-new/base_template/src/app.rs.t),你只需把 Worker 的实例registerQueue

#[async_trait] impl Hooks for App { //.. async fn connect_workers(ctx: &AppContext, queue: &Queue) -> Result<()> { queue.register(DownloadWorker::build(ctx)).await?; Ok(()) } // .. }

register(src/bgworker/mod.rs)内部以W::class_name()为键,把 Worker 实例登记到对应后端(Redis/Postgres/SQLite)的 JobRegistry 中,之后 Worker 进程启动时就能根据任务标识找到正确的处理函数。

用脚手架一键生成 Worker

使用loco generate子命令自动生成 Worker 文件与测试模板:

cargo loco generate worker report_worker

生成器(模板见 loco-gen/src/templates/worker/worker.t 与 loco-gen/src/templates/worker/test.t)会做三件事:

  1. src/workers/<name>.rs生成Worker结构与WorkerArgs参数结构,实现buildclass_nametagsperform
  2. src/workers/mod.rs追加pub mod <name>;
  3. src/app.rsconnect_workers后自动注入queue.register(crate::workers::<name>::Worker::build(ctx)).await?;注册语句;
  4. tests/workers/<name>.rs生成对应的测试骨架。

命令行入口定义见 src/cli.rs 的Generate Worker分支。

BackgroundWorker trait 全解析与 class_name() 命名规则

BackgroundWorker<A>是定义后台任务的核心接口(src/bgworker/mod.rs),其中A是任务参数类型,需要满足Send + Sync + Serialize + 'static。它提供的成员方法:

  • build(ctx: &AppContext) -> Self:使用应用上下文创建 Worker 实例,注册与执行时都会调用;
  • perform(&self, args: A) -> Result<()>:执行任务逻辑的主方法,接收反序列化后的参数;
  • queue() -> Option<String>:可选,指定自定义队列名(默认None)。若队列后端支持自定义/优先级队列,可在这里返回队列名;
  • tags() -> Vec<String>:可选,声明任务标签(默认空向量);
  • class_name() -> String:返回 Worker 的类名标识,自动从结构体名推导
  • perform_later(ctx: &AppContext, args: A) -> Result<()>:静态方法,把任务入队以便稍后执行。

class_name() 的推导规则

class_name()用于在任务队列中唯一标识你的 Worker。默认实现(src/bgworker/mod.rs)分三步:

  1. 取结构体名(如DownloadWorker);
  2. 剥离模块路径(如my_app::workers::DownloadWorker只保留最后的DownloadWorker);
  3. 转换为 UpperCamelCase 格式。

由于任务入队时需要字符串标识来匹配对应的处理 Worker,class_name()自动生成该标识;若你需要自定义命名方案,可以覆盖这个方法。默认实现基于std::any::type_name::<Self>()取末段再经heck::ToUpperCamelCase转换:

// Example of how class_name works: struct download_worker; impl BackgroundWorker<Args> for download_worker { // class_name() would return "DownloadWorker" // No need to override this unless you need custom naming }

在 Worker 中使用共享状态

Worker 需要全局共享状态时,可以参考 全局应用级状态 的做法,一般用lazy_static建立单一共享状态,然后在 Worker 里直接引用(build时把AppContext克隆进 Worker 即可访问其中的资源)。

一个更推荐的原则:如果状态可序列化,强烈建议通过WorkerArgs传参,而不是从外部全局状态读取——这能保证任务在跨进程、跨机器的队列模式下也能拿到完整数据,避免依赖进程内的全局变量。

测试 Worker:ForegroundBlocking 模式下的同步验证

Loco 让 Worker 测试非常直接:把workers.mode设为ForegroundBlocking,任务会同步阻塞执行,测试等待任务完成后即可断言其副作用。建议把测试统一放在tests/workers/目录集中管理;也可以直接使用 worker 生成器自动创建的测试,免去手工配置的麻烦。

测试模板示例(完整版见 loco-gen/src/templates/worker/test.t):

use loco_rs::testing::prelude::*; #[tokio::test] #[serial] async fn test_run_report_worker_worker() { // Set up the test environment let boot = boot_test::<App, Migrator>().await.unwrap(); // Execute the worker in 'ForegroundBlocking' mode, preventing it from running asynchronously assert!( ReportWorkerWorker::perform_later(&boot.app_context, ReportWorkerWorkerArgs {}) .await .is_ok() ); // Include additional assert validations after the execution of the worker }

要点:#[tokio::test]提供异步运行时,#[serial](来自serial_test)确保测试串行执行;boot_test完成应用引导;由于测试环境配置为ForegroundBlockingperform_later内部直接同步执行perform,因此assert!(...is_ok())通过即代表任务成功跑完。

通过 CLI 管理任务队列:cancel / tidy / purge / dump / import / requeue

任务队列的 CLI 管理能力提供了对任务生命周期的精细控制:取消、清理、移除过期任务、导出任务详情与导入任务,保证队列高效有序。入口是jobs子命令(定义见 src/cli.rs 与 src/cli.rs,具体实现见 src/cli.rs 的handle_job_command):

Managing jobs queue Usage: demo_app-cli jobs [OPTIONS] <COMMAND> Commands: cancel Cancels jobs with the specified names, setting their status to `cancelled` tidy Deletes jobs that are either completed or cancelled purge Deletes jobs based on their age in days dump Saves the details of all jobs to files in the specified folder import Imports jobs from a file requeue Change `processing` jobs older than the specified maximum age (in minutes) back to `queued` help Print this message or the help of the given subcommand(s) Options: -e, --environment <ENVIRONMENT> Specify the environment [default: development] -h, --help Print help -V, --version Print version

各子命令的能力与底层实现如下:

  • cancel:按名称取消指定任务,将其状态更新为cancelled。适用于停止不再需要、已过时或发现 Bug 后希望阻止其继续处理的任务。底层调用Queue::cancel_jobs(src/bgworker/mod.rs)。
  • tidy:删除已完成(completed)或已取消(cancelled)的任务,保持队列干净高效。底层调用Queue::clear_by_status(src/bgworker/mod.rs)。
  • purge:按任务年龄(天)删除过期任务,适合精简队列。可选--max-age(默认 90 天)、--status(按状态过滤)、--dump(删除前先导出到文件)。底层调用Queue::clear_jobs_older_than(src/bgworker/mod.rs)。
  • dump:把所有任务详情导出到指定目录下的 YAML 文件。实用技巧:可以先用--dump导出任务,手工修改文件中的任务参数,再用import把修改后的任务重新导回系统。底层调用Queue::dump(src/bgworker/mod.rs),文件命名为loco-dump-jobs-<时间戳>.yaml
  • import:从外部 YAML 文件导入任务,便于恢复或批量新增任务。底层调用Queue::import(src/bgworker/mod.rs)。
  • requeue:把processing状态超过指定分钟数(--from-age,默认 0 分钟)的“僵死”任务重新置回queued,用于故障恢复。底层调用Queue::requeue(src/bgworker/mod.rs)。

任务状态枚举JobStatus(src/bgworker/mod.rs)包含五种:queuedprocessingcompletedfailedcancelled,贯穿整个任务生命周期。上述能力在仓库中均有测试覆盖,例如 src/bgworker/mod.rs 中的can_dump_jobscat_import_jobs_form_file(后者使用 tests/fixtures/queue/jobs.yaml 作为导入样本),可作参考。

另外,社区提供的 Loco admin job 项目还提供了可视化的任务队列管理界面(UI 方式),适合希望以图形化方式查看与操作队列的场景。

小结:按场景选择 Worker 方案

场景推荐方案
单机、无需额外基础设施BackgroundAsync(默认 async 配置),同进程 Tokio 异步执行
需要把负载分发到多台服务器BackgroundQueue+ Redis/Postgres/SQLite,配合cargo loco start --worker独立进程
需要同步执行、便于测试ForegroundBlocking,任务阻塞执行,测试可直接断言结果
资源密集型任务隔离利用标签过滤,启动专职 Worker 只消费特定标签
队列运维与故障恢复jobs子命令:cancel / tidy / purge / dump / import / requeue

核心结论:Loco 的 Worker 体系以BackgroundWorkertrait 为统一契约,通过一行 YAML 配置切换执行模式与队列后端,业务代码始终只写perform_later,既保留了 Rust 强类型参数的安全感,又获得了类似 ActiveJob 的部署灵活性。相关源码入口:任务抽象与实现位于 src/bgworker/mod.rs、配置结构位于 src/config.rs、CLI 命令位于 src/cli.rs,Worker 脚手架模板位于 loco-gen/src/templates/worker/,新项目自带示例见 loco-new/base_template/src/workers/downloader.rs。

【免费下载链接】loco🚂 🦀 The one-person framework for Rust for side-projects and startups项目地址: https://gitcode.com/GitHub_Trending/lo/loco

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询