外观
任务调度
概述
后端提供两套互补的调度能力,分别面向「耗时任务的异步执行」与「按时间周期触发」两种场景:
- 异步任务队列(task_queue):基于 arq(Redis 之上的异步任务队列),把耗时/重计算任务(如 AI 文档解析、向量化)从请求链路中剥离,异步消费,支持水平扩容。
- 定时任务(cron):基于 APScheduler 的
AsyncIOScheduler,配合framework/common/cron.py的 Quartz 风格表达式校验与转换,按 Cron 周期驱动处理器执行,执行记录落库可追溯。
两者都依赖 Redis,但用途不同:任务队列用 Redis 作为消息通道,定时调度用 Redis 作为任务持久化存储(JobStore)。
异步任务队列(task_queue)
框架对 arq 做了薄封装,业务层只感知 enqueue(函数名, 业务主键) 语义,不直接依赖 arq 类型。设计上的关键约定:队列消息仅携带业务主键(如 task_id),完整上下文以数据库为准,从而保证任务可重入、幂等,并且 worker 崩溃重启后能以 DB 状态为权威源继续推进。
核心文件位于 app/framework/task_queue/:
client.py:TaskQueueClient惰性单例连接池,提供enqueue/close。settings.py:TaskQueueSettings,复用全局 Redis 配置,定义队列名、并发、超时等行为参数。
工作流程
以「AI 知识库文档解析」为例,完整的投递—消费链路如下:
投递任务
业务侧通过 TaskQueueClient.enqueue 投递,函数名需与 worker 端注册的函数一致:
python
from app.framework.task_queue.client import TaskQueueClient
# 函数名对应 WorkerSettings.functions 中注册的可调用对象
# payload 仅放业务主键,复杂上下文在执行端从数据库读取
await TaskQueueClient.enqueue(
"ai_parse_document",
{"task_id": task.id},
job_id=f"ai_parse_task_{task.id}", # 去重 ID,同 ID 任务未执行完时重复投递会被忽略
)enqueue 返回值语义:
True:投递成功;False:因相同job_id任务已在队列中而被忽略(幂等保护)。
job_id 是队列层的幂等键,建议由「业务类型 + 业务主键」拼成,避免重复任务堆积。
编写并注册 Worker
Worker 是独立进程,生产环境可多实例启动,由 Redis 队列自动分摊消费:
bash
uv run arq app.api.v1.module_ai.framework.task_queue.worker.WorkerSettingsWorker 入口通过 WorkerSettings 声明消费的函数与行为,队列名必须与生产端 enqueue 使用的队列一致,否则任务不会被消费:
python
from arq.connections import RedisSettings
from app.api.v1.module_ai.knowledge.document.parse_worker import ai_parse_document
from app.framework.task_queue.settings import task_queue_settings
class WorkerSettings:
functions = [ai_parse_document] # 注册可消费的任务函数
queue_name = task_queue_settings.TASK_QUEUE_NAME
redis_settings = task_queue_settings.redis_settings
max_jobs = task_queue_settings.TASK_QUEUE_MAX_JOBS
job_timeout = task_queue_settings.TASK_QUEUE_JOB_TIMEOUT
max_tries = task_queue_settings.TASK_QUEUE_MAX_TRIES
health_check_interval = task_queue_settings.TASK_QUEUE_HEALTH_CHECK_INTERVAL执行端函数签名遵循 arq 约定:第一个参数为 ctx(含 Redis 连接),其余参数与 enqueue 的 payload 对应:
python
async def ai_parse_document(ctx: dict, *, task_id: int) -> None:
redis = ctx.get("redis")
# 1. 按 task_id 从 DB 加载任务与文档
# 2. 分阶段执行,并用回调实时回写进度 / 推送 WebSocket
# 3. 阶段间检查取消标志,支持手动中止
# 4. 异常时回写 FAILED 状态,成功时回写 COMPLETED配置项
TaskQueueSettings 的关键配置(可按环境在 .env 覆盖):
TASK_QUEUE_NAME:队列名,默认ai_parse;TASK_QUEUE_MAX_JOBS:单 worker 并发执行的任务数,默认2;TASK_QUEUE_JOB_TIMEOUT:单任务超时(秒),默认3600;TASK_QUEUE_MAX_TRIES:队列层最大重试次数,默认1(业务失败重试建议由任务表的retry_count承担,而非依赖此处自动重试);TASK_QUEUE_HEALTH_CHECK_INTERVAL:健康检查间隔(秒),默认30。
Redis 连接复用全局 settings 中的 Redis 配置,通过 DSN 自动构造 RedisSettings,无需单独维护连接串。
定时任务(cron)
定时任务用于按周期自动执行处理器(如数据清理、报表生成、状态同步等),实现位于 app/api/v1/module_infra/job/,由框架层 framework/common/cron.py 提供表达式校验能力。
调度器设计
调度器是一个独立的 AsyncIOScheduler,使用独立的 RedisJobStore 持久化任务定义,执行器为 AsyncIOExecutor(原生驱动协程),时区固定为 Asia/Shanghai:
python
infra_scheduler = AsyncIOScheduler(
jobstores={"default": RedisJobStore(..., jobs_key="infra_apscheduler_jobs", run_times_key="infra_apscheduler_run_times")},
executors={"default": AsyncIOExecutor()},
job_defaults={"coalesce": True, "max_instances": 1},
timezone="Asia/Shanghai",
)默认 coalesce=True 且 max_instances=1,保证同一任务错过多次触发时合并为一次执行,避免堆积。
工作流程
Cron 表达式
表达式采用 Quartz 风格,支持 6 位(秒 分 时 日 月 周) 或 7 位(追加 年),由 framework/common/cron.py 的 CronUtil 校验合法性。框架在构造 APScheduler 触发器时会做等价转换:
- Quartz 的
?(“不指定”)等价替换为*; - 日期字段的
L转换为 APScheduler 的last(当月最后一天); - 其余字段(
* - , / W #等)原样透传。
若接入新的定时任务,建议在后端校验表达式,避免非法配置进入调度器:
python
from app.framework.common.cron import CronUtil
if not CronUtil.validate_cron_expression("0 0 2 * * ?"):
raise ValueError("无效的 Cron 表达式")处理器(handler)机制
调度器实际执行的任务体是 run_infra_job(job_id):它按任务 ID 从 infra_job 表加载定义,再动态导入并执行 handler_name 指向的处理器。处理器通过 handler_name + handler_param 两个字段声明:
handler_name:可调用对象路径,形如module_infra.file.task.say_hello;若使用:分隔则按「模块:属性」解析,否则按「模块.属性」解析。缺少app.api.v1.前缀时会自动补全后重试(兼容历史数据)。handler_param:JSON 对象,作为关键字参数传给处理器。
处理器可以是同步函数或异步协程,调度器会自动 await 协程结果:
python
# app/api/v1/module_infra/file/task.py
async def say_hello(name: str = "world"):
return f"hello {name}"json
// 任务配置示例
{ "handler_name": "module_infra.file.task.say_hello", "handler_param": { "name": "mars" } }执行与日志
每次触发都会执行 run_infra_job,其执行结果(成功/失败、返回数据或错误信息、耗时、执行序号)统一写入 infra_job_log 表,前端可在「任务日志」中查看,实现执行可追溯。run_infra_job 为协程,既被调度器到点触发,也可被「立即执行一次」接口直接调用(手动触发仅记录一次日志,不依赖 Cron)。
生命周期与启动
调度器的启动与关闭挂载在 FastAPI 的 lifespan 中(app_init.py):
- 启动时:若未运行则
start(),随后remove_all_jobs()清空,再调用sync_jobs()把数据库中**状态为启用(NORMAL)**的任务全量重新注册,返回(成功数, 失败数)。 - 关闭时:
shutdown(wait=False)优雅停止。
因此新增/修改/启停任务后,除了通过接口即时生效外,也可调用「同步」接口触发 sync_jobs() 全量重载,保证调度器与数据库一致。
管理端与前端集成
定时任务的管理界面与接口集中在 基础设施模块(module_infra,管理端文档见 05),后端手册此处只说明能力边界与对接方式。
后端为前端提供的能力:
- 查询/新增/修改/删除定时任务,启用即注册到调度器,暂停即从调度器移除;
- 「立即执行一次」:绕过 Cron 直接调用
run_infra_job,便于调试; - 「获取下次执行时间」:基于 Cron 表达式计算未来若干次触发点,前端可展示「预计下次运行」;
- 「同步」:全量重载启用任务到调度器;
- 「任务日志」:分页查看每次执行的开始/结束时间、耗时、状态与错误信息,支持导出 Excel。
前端建议的展示要点:
- 任务列表展示
cron_expression与next_times(下次执行时间),让用户直观确认调度是否生效; - 进入任务详情后切换「执行日志」标签页,按
execute_index倒序展示,红色高亮FAILURE记录并展示error字段; - 提供「立即执行」按钮(触发
trigger)与「状态切换」开关(启停),操作后刷新列表状态。
异步任务队列(如文档解析)通常不直接暴露给前端管理,而是通过业务接口(上传/解析/取消)间接使用;前端主要消费其进度:后端在执行端分阶段通过 WebSocket 推送解析进度,前端据此渲染进度条,轮询接口仅作为兜底。