Skip to content

任务调度 ​

概述 ​

后端提供两套互补的调度能力,分别面向「耗时任务的异步执行」与「按时间周期触发」两种场景:

  • 异步任务队列(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.WorkerSettings

Worker 入口通过 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 推送解析进度,前端据此渲染进度条,轮询接口仅作为兜底。