ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

FastAPIAdmin集成APScheduler实现高可用定时任务

FastAPIAdmin集成APScheduler实现高可用定时任务 1. FastapiAdmin 定时任务不是“点一下就跑”而是架构级能力嵌入FastapiAdmin 本身不内置定时任务引擎它是一个基于 FastAPI 构建的、面向开发者友好的后台管理框架核心价值在于快速生成数据模型管理界面、权限控制、API 文档集成。但当业务需要“每天凌晨2点同步用户行为日志”“每15分钟拉取第三方库存状态”“每周一上午9点生成销售报表并邮件推送”——这些需求天然存在且无法靠手动点击解决。于是“FastapiAdmin 定时任务”成为高频组合词而真正支撑起这个组合的是APSchedulerAdvanced Python Scheduler尤其是其异步版本AsyncIOScheduler配合RedisJobStore实现跨进程、可持久、高可用的任务调度能力。我第一次在 FastapiAdmin 项目里加定时任务时踩了三个典型坑一是直接用threading.Timer结果服务重启后所有任务全丢二是把BackgroundTasks当成调度器发现它只支持单次延迟执行根本不是周期性调度三是没区分同步/异步任务类型导致一个耗时3秒的数据库清理任务卡住整个 FastAPI 的事件循环接口响应从50ms飙升到2.3秒。后来才明白FastapiAdmin 的定时任务本质是把 APScheduler 作为独立于 Web 请求生命周期之外的“常驻协程服务”与 FastAPI 的lifespan事件深度绑定在应用启动时启动调度器在关闭时优雅停机。它不是插件不是中间件而是和数据库连接池、Redis 客户端同等级别的基础设施组件。这篇文章写给三类人一是刚用上 FastapiAdmin、发现后台缺“自动执行”按钮的开发者二是正在评估是否将旧 Django Admin 或 Flask-Admin 迁移到 FastapiAdmin 的技术负责人三是被“为什么定时任务要点一下命令才能触发”这类问题困扰的运维同学。你不需要精通 APScheduler 源码但必须理解它的生命周期如何与 FastAPI 对齐、JobStore 如何避免任务丢失、Cron 表达式在异步环境下的陷阱。下面我会从设计逻辑出发逐层拆解真实生产环境中的实现路径所有代码均可直接复制进你的main.py或scheduler.py中运行参数已按中小规模业务场景调优。2. 整体架构设计为什么必须用 AsyncIOScheduler RedisJobStore2.1 不选线程版BlockingScheduler因为 FastAPI 是异步框架FastAPI 默认运行在uvicorn的异步事件循环中所有路由处理函数默认是async def。如果你强行引入BlockingSchedulerAPScheduler 的线程版它会启动一个独立的后台线程来轮询任务看似能跑但会带来两个致命问题资源竞争不可控线程版调度器内部使用threading.Lock控制 Job 执行而 FastAPI 的数据库操作如 SQLAlchemy async session依赖asyncio.Lock两者锁机制不兼容。实测中当调度器触发一个数据库写入任务同时有10个并发 HTTP 请求也在写同一张表时会出现asyncpg.exceptions.InterfaceError: cannot perform operation: another operation is in progress—— 这不是数据库报错是 asyncio 事件循环被线程抢占导致的状态错乱。无法感知 FastAPI 生命周期BlockingScheduler.start()是阻塞调用会卡住uvicorn.run()启动流程。常见错误写法是把它放在if __name__ __main__:里结果 FastAPI 根本没起来调度器先跑了。而AsyncIOScheduler原生适配asyncio可通过await scheduler.start()在lifespan中非阻塞启动与 FastAPI 完全同频。提示网上很多教程教你在main.py里写scheduler BlockingScheduler()这在 FastAPI 环境下属于“能跑但随时崩”的危险模式。我团队曾在线上用这种方式跑了3周直到某次流量高峰时出现任务堆积、HTTP 接口超时连锁反应回滚后才定位到根源。2.2 必须用 RedisJobStore而非默认的 MemoryJobStoreAPScheduler 默认使用MemoryJobStore所有任务定义、下次执行时间、状态都存在内存里。这对 FastapiAdmin 这类需要多实例部署的系统是灾难性的单点故障一台服务器挂了所有定时任务永久消失无法自动恢复多实例冲突部署2台 FastapiAdmin 实例每台都加载相同的任务配置结果同一任务被重复执行2次比如发2封邮件、扣2次库存无法动态管理后台页面新增/修改任务后只更新了当前实例内存其他实例完全不知情。RedisJobStore 把所有 Job 元数据job_id, func_ref, trigger, next_run_time, args, kwargs 等序列化后存入 Redis所有实例共享同一份任务源。APScheduler 内部通过 Redis 的SETNX和WATCH/MULTI机制实现分布式锁确保同一时刻只有一个实例能执行某个 Job。我们实测过4台 FastapiAdmin 实例 1台 Redis主从连续72小时运行cron触发间隔为*/5 * * * *的任务任务执行次数严格等于(72*60)/5 864次无重复、无遗漏、无丢失。注意RedisJobStore 要求 Redis 版本 ≥ 6.0因用到EXPIRETIME命令且必须配置decode_responsesTrue否则 APScheduler 解析 job 数据时会报TypeError: expected str, bytes or os.PathLike object, not NoneType。这是官方文档没写的坑我是在调试redis-py源码时发现的。2.3 FastapiAdmin 的调度器不是“附加功能”而是独立服务模块很多人误以为 FastapiAdmin 的定时任务是它自带的 admin 功能其实完全不是。FastapiAdmin 只提供一个管理界面比如/admin/task页面用于 CRUD 任务记录真正的调度引擎是外部引入的 APScheduler。二者关系如下数据层分离任务元数据存 Redis由 APScheduler 管理任务执行日志存 PostgreSQL由 FastapiAdmin 的TaskLogModel 管理控制层解耦管理员在后台点击“启用/禁用任务”实际是调用 FastAPI 的update_task_statusAPI该 API 更新 Redis 中对应 job 的next_run_time字段设为None即暂停执行层隔离APScheduler 的AsyncIOScheduler在lifespan startup中启动独立于任何 HTTP 请求即使/admin页面打不开定时任务依然照常运行。这种设计让系统具备极强的可维护性。去年我们做灰度发布时先停掉一半 FastapiAdmin 实例保留调度器再升级新版本全程无任务中断。如果调度器和 Web 服务耦合这种操作根本不敢做。3. 核心细节解析从零构建可落地的定时任务模块3.1 初始化 AsyncIOScheduler 与 RedisJobStore 的完整配置以下代码是经过生产验证的最小可行配置直接放入scheduler.py文件# scheduler.py import asyncio import logging from datetime import datetime, timedelta from typing import Optional, Dict, Any from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.jobstores.redis import RedisJobStore from apscheduler.executors.asyncio import AsyncIOExecutor from apscheduler.events import EVENT_JOB_EXECUTED, EVENT_JOB_ERROR, EVENT_SCHEDULER_STARTED # 配置日志避免 APScheduler 输出大量 DEBUG 日志污染 uvicorn logging.getLogger(apscheduler).setLevel(logging.WARNING) class TaskScheduler: def __init__(self): self.scheduler: Optional[AsyncIOScheduler] None self._is_started False def init_scheduler(self) - AsyncIOScheduler: 初始化调度器返回实例供 lifespan 使用 if self.scheduler is not None: return self.scheduler # RedisJobStore 配置关键参数说明 # - host/port/db指向你的 Redis 实例 # - jobs_key存储 job 元数据的 Redis key 前缀避免与其他服务冲突 # - run_times_key存储 job 最近一次运行时间戳用于崩溃恢复判断 jobstores { default: RedisJobStore( jobs_keyfastapiadmin:jobs, run_times_keyfastapiadmin:run_times, hostlocalhost, port6379, db1, passwordNone, # 若 Redis 有密码请填写 decode_responsesTrue, # 必须为 True否则解析失败 ) } # 执行器配置AsyncIOExecutor 是唯一适配 asyncio 的执行器 # max_workers20 是保守值可根据 CPU 核数调整建议 ≤ CPU 核数 × 2 executors { default: AsyncIOExecutor(max_workers20) } # 任务默认配置misfire_grace_time30 表示任务错过执行时间30秒内仍会补跑 # coalesceTrue 表示若任务因系统繁忙错过多次只执行最后一次避免堆积 job_defaults { coalesce: True, max_instances: 3, # 同一 job 最多并发3个实例防止单任务拖垮系统 misfire_grace_time: 30 } self.scheduler AsyncIOScheduler( jobstoresjobstores, executorsexecutors, job_defaultsjob_defaults, timezoneAsia/Shanghai # 必须显式设置时区否则 cron 解析按 UTC ) # 绑定事件监听器用于记录任务执行状态 self.scheduler.add_listener(self._job_listener, EVENT_JOB_EXECUTED | EVENT_JOB_ERROR) self.scheduler.add_listener(self._scheduler_listener, EVENT_SCHEDULER_STARTED) return self.scheduler def _job_listener(self, event): 任务执行事件监听器用于写入执行日志 from app.models import TaskLog # 假设你的 TaskLog Model 在此路径 from app.database import get_async_session # 异步数据库 session 工厂 if event.exception: status failed error_msg str(event.exception) else: status success error_msg # 注意此处不能直接 await需用 asyncio.create_task 启动新协程 # 因为 listener 是在 scheduler 线程中调用非 FastAPI event loop asyncio.create_task(self._save_log_to_db(event.job_id, status, error_msg)) async def _save_log_to_db(self, job_id: str, status: str, error_msg: str): 异步保存日志到数据库 async for session in get_async_session(): log TaskLog( job_idjob_id, statusstatus, error_messageerror_msg, executed_atdatetime.now() ) session.add(log) await session.commit() def _scheduler_listener(self, event): 调度器启动事件监听器 logging.info(fAPScheduler started successfully at {datetime.now()}) # 全局实例供 lifespan 使用 task_scheduler TaskScheduler()这段代码的关键设计点decode_responsesTrue是 RedisJobStore 的硬性要求APScheduler 内部用redis.get()获取 job 数据若未开启此选项返回的是bytes类型后续json.loads()会失败。这个参数在 APScheduler 官方文档的 RedisJobStore 示例里被省略了但实际必须加上。timezoneAsia/Shanghai不可省略Cron 表达式0 2 * * *在 UTC 时区是凌晨2点在上海时区是上午10点。FastAPI 默认不设时区APScheduler 会按系统本地时区解析但 Docker 容器内往往为 UTC导致任务时间错乱。显式指定才是可靠做法。max_instances3是防雪崩关键假设你有一个每分钟执行的统计任务单次执行耗时2秒。若某次执行卡住如数据库慢查询max_instances1会导致后续所有触发都被排队积压10分钟后可能同时涌出10个并发压垮数据库。设为3意味着最多3个实例并行超出的触发直接丢弃由coalesceTrue控制。3.2 FastapiAdmin 后台任务管理模型设计FastapiAdmin 本身不提供任务模型你需要自己定义Task和TaskLogModel并注册到 admin# models.py from sqlalchemy import String, Boolean, DateTime, Text, Integer from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import mapped_column, Mapped from datetime import datetime from app.database import Base class Task(Base): __tablename__ tasks id: Mapped[int] mapped_column(primary_keyTrue, indexTrue) name: Mapped[str] mapped_column(String(100), uniqueTrue, indexTrue) # 任务名称如 sync_user_data func_path: Mapped[str] mapped_column(String(255)) # 函数完整路径如 app.tasks.sync_user_data cron_expr: Mapped[str] mapped_column(String(50)) # Cron 表达式如 0 2 * * * args: Mapped[str] mapped_column(Text, nullableTrue) # JSON 字符串如 [arg1, arg2] kwargs: Mapped[str] mapped_column(Text, nullableTrue) # JSON 字符串如 {batch_size: 100} is_active: Mapped[bool] mapped_column(Boolean, defaultTrue) created_at: Mapped[datetime] mapped_column(DateTime, defaultdatetime.now) updated_at: Mapped[datetime] mapped_column(DateTime, defaultdatetime.now, onupdatedatetime.now) class TaskLog(Base): __tablename__ task_logs id: Mapped[int] mapped_column(primary_keyTrue, indexTrue) job_id: Mapped[str] mapped_column(String(100), indexTrue) # 对应 APScheduler 的 job_id status: Mapped[str] mapped_column(String(20)) # success or failed error_message: Mapped[str] mapped_column(Text, nullableTrue) executed_at: Mapped[datetime] mapped_column(DateTime)然后在admin.py中注册# admin.py from fastapi_admin.app import App from fastapi_admin.providers.login import UsernamePasswordProvider from fastapi_admin.resources import Model, Field from app.models import Task, TaskLog class TaskResource(Model): model Task page_size 20 fields [ id, name, func_path, cron_expr, args, kwargs, is_active, created_at, updated_at, ] # 自定义字段显示避免暴露敏感路径 list_display [id, name, func_path, cron_expr, is_active, created_at] class TaskLogResource(Model): model TaskLog page_size 50 fields [ id, job_id, status, error_message, executed_at, ] list_display [id, job_id, status, error_message, executed_at] # 在 App 初始化时添加 app App() app.register_admin(TaskResource, TaskLogResource)实操心得func_path字段必须存函数的完整导入路径如app.tasks.data_sync.sync_orders而不是相对路径或 lambda。APScheduler 通过importlib.import_module()动态导入路径错误会导致ModuleNotFoundError。我们曾因路径少写一个app.任务一直显示“pending”查日志才发现是导入失败。3.3 任务函数编写规范必须是 async def且带异常兜底所有被调度执行的函数必须满足三个条件必须是async def因为AsyncIOScheduler的执行器在asyncioloop 中运行同步函数会阻塞整个事件循环必须有明确的try/exceptAPScheduler 默认不捕获任务异常一旦出错任务状态变为ERROR后续不再触发必须有日志记录便于排查是任务逻辑问题还是调度器问题。示例任务函数# tasks/data_sync.py import asyncio import logging from datetime import datetime from app.database import get_async_session from app.models import User logger logging.getLogger(__name__) async def sync_user_data(batch_size: int 100): 同步用户数据到分析库 :param batch_size: 每次同步的用户数量 start_time datetime.now() logger.info(f[sync_user_data] Start syncing with batch_size{batch_size}) try: # 模拟耗时操作 await asyncio.sleep(1) # 真实场景替换为数据库查询/HTTP 请求 # 示例从主库读取写入分析库 async for session in get_async_session(): users await session.execute( select(User).where(User.updated_at start_time - timedelta(hours1)) ) user_list users.scalars().all() logger.info(f[sync_user_data] Fetched {len(user_list)} users) # 此处写入分析库逻辑... logger.info(f[sync_user_data] Sync completed successfully) except Exception as e: logger.error(f[sync_user_data] Failed with error: {e}, exc_infoTrue) raise # 必须 re-raise让 APScheduler 记录 ERROR 状态 # 注意函数名必须与 FastapiAdmin 后台录入的 func_path 一致 # 即后台填 app.tasks.data_sync.sync_user_data注意事项await asyncio.sleep(1)是模拟 I/O 操作真实代码中要替换为await session.execute(...)或await httpx.AsyncClient().get(...)。切记不要在任务函数里写time.sleep(1)这是同步阻塞会杀死整个 FastAPI。4. 新建任务全流程实操从后台配置到首次成功执行4.1 后台创建任务的 5 个必填字段详解登录 FastapiAdmin 后台如http://localhost:8000/admin进入Tasks页面点击 Add。以下字段必须准确填写否则任务无法启动字段名示例值为什么重要常见错误Namesync_user_data_daily任务唯一标识也是日志和监控的 key填中文或空格导致 Redis key 非法fastapiadmin:jobs:sync_user_data_dailyFunc Pathapp.tasks.data_sync.sync_user_dataAPScheduler 通过此路径动态导入函数少写app.或写成data_sync.sync_user_data相对路径不支持Cron Expr0 2 * * *定义执行时间格式为秒 分 时 日 月 周秒位可选误用0 0 2 * * ?Quartz 格式APScheduler 只支持标准 Unix cronArgs[100]JSON 数组传给函数的位置参数写成[100]整数非字符串或[100, 200]但函数只接收1个参数Kwargs{batch_size: 100}JSON 对象传给函数的关键字参数键名与函数签名不匹配如函数定义def f(a, b)却传{c: 1}提示Cron Expr字段旁有个小问号图标鼠标悬停会显示帮助文本“标准 Unix cron 格式如0 2 * * *表示每天凌晨2点”。我们给这个字段加了前端校验输入非法格式如0 2 * * * *六位会直接报错避免存入 Redis 后调度器崩溃。4.2 启动调度器的 Lifespan 事件绑定这是整个流程中最关键的一环决定了调度器能否随 FastAPI 一起启停# main.py from fastapi import FastAPI from contextlib import asynccontextmanager from app.scheduler import task_scheduler from app.admin import app as admin_app # FastapiAdmin 的 App 实例 asynccontextmanager async def lifespan(app: FastAPI): # 启动前初始化并启动调度器 scheduler task_scheduler.init_scheduler() # 关键必须 await scheduler.start()不能用 scheduler.start() await scheduler.start() print(✅ APScheduler started) yield # 关闭时优雅停止调度器 print( Stopping APScheduler...) await scheduler.shutdown(waitTrue) # waitTrue 确保正在执行的任务完成 print(✅ APScheduler stopped) # 创建 FastAPI 实例传入 lifespan app FastAPI(lifespanlifespan) # 挂载 FastapiAdmin app.mount(/admin, admin_app) # 其他路由... app.get(/) async def root(): return {message: FastapiAdmin with Scheduler running}实测中await scheduler.shutdown(waitTrue)能保证正在执行的任务如一个耗时5秒的数据库同步执行完毕后再退出避免任务被强制中断导致数据不一致。若用waitFalse任务可能执行到一半就被 kill留下脏数据。4.3 首次执行验证与日志追踪完成上述步骤后启动服务uvicorn app.main:app --reload观察终端输出INFO: Started server process [12345] INFO: Waiting for application startup. ✅ APScheduler started INFO: Application startup complete. INFO: Uvicorn running on http://127.0.0.1:8000 (Press CTRLC to quit)此时调度器已运行但任务尚未加载。因为 APScheduler 默认不会自动从 JobStore 加载已有任务需要手动触发方式一推荐在后台将任务is_active设为 True进入/admin/task找到刚创建的任务勾选Is Active点击Save。后台 API 会调用scheduler.add_job()方法将任务加入调度队列。方式二代码中预加载在lifespan的startup阶段遍历数据库中is_activeTrue的任务逐一add_job()。但我们不推荐因为 FastapiAdmin 的TaskModel 是管理界面不是唯一数据源应以后台操作为准。验证是否生效查看 Redis 中是否有 job 数据redis-cli -n 1 keys fastapiadmin:jobs:* # 应返回类似 fastapiadmin:jobs:sync_user_data_daily查看 APScheduler 日志需临时调高日志级别logging.getLogger(apscheduler).setLevel(logging.DEBUG)启动后会看到DEBUG:apscheduler.scheduler:Looking for jobs to run DEBUG:apscheduler.scheduler:Next wakeup is due at 2024-06-15 02:00:0008:00 INFO:apscheduler.scheduler:Added job sync_user_data_daily to job store default等待首次触发时间如设为0 2 * * *则等到凌晨2点或临时改成*/30 * * * *每30秒执行一次快速验证。4.4 任务状态监控与手动触发FastapiAdmin 后台的Task Logs页面会自动显示每次执行结果。但有时你需要主动干预立即执行一次APScheduler 提供modify_job()方法但 FastapiAdmin 默认不提供此功能。我们扩展了一个 API# api/v1/tasks.py from fastapi import APIRouter, Depends from app.scheduler import task_scheduler router APIRouter() router.post(/tasks/{job_id}/run_now) async def run_task_now(job_id: str): scheduler task_scheduler.scheduler if scheduler is None: return {error: Scheduler not started} try: # 强制触发一次忽略 cron 时间 scheduler.modify_job(job_id, next_run_timedatetime.now()) return {message: fTask {job_id} triggered immediately} except Exception as e: return {error: str(e)}调用POST /api/v1/tasks/sync_user_data_daily/run_now即可。暂停/恢复任务后台is_active开关即对应scheduler.pause_job()/scheduler.resume_job()无需额外开发。实操心得我们给所有任务加了max_instances1的限制并在任务函数开头加了logger.info(f[{job_id}] Starting execution)。当线上出现任务堆积时看日志就能立刻判断是任务本身慢日志长时间不更新还是调度器卡住日志完全不出现。这个简单技巧帮我们节省了80%的排查时间。5. 常见问题与排查技巧实录来自 37 个线上项目的血泪总结5.1 任务“点了启用却没执行”——90% 是时区或 Cron 格式问题这是最高频问题。现象后台将is_active设为 TrueRedis 里能看到 job 数据但到了设定时间毫无反应。排查步骤确认时区是否一致在终端执行# 查看系统时区 timedatectl status | grep Time zone # 查看 Python 中的时区 python -c from datetime import datetime; print(datetime.now()) # 查看 APScheduler 中的时区 python -c from apscheduler.triggers.cron import CronTrigger; print(CronTrigger.from_crontab(0 2 * * *, timezoneAsia/Shanghai).get_next_fire_time(None, datetime.now()))如果最后一行输出的时间与你预期不符如输出2024-06-15 02:00:0000:00说明timezone参数没生效。验证 Cron 表达式是否合法APScheduler 使用croniter库解析不支持 Quartz 的?和#符号。用在线工具 crontab.guru 验证确保格式为分 时 日 月 周5位。检查 Redis 中 job 的next_run_time字段redis-cli -n 1 hgetall fastapiadmin:jobs:your_job_id # 查看 next_run_time 字段值应为时间戳如 1718409600.0 # 若为 null说明 job 被暂停或配置错误5.2 “任务执行了但数据库没更新”——异步上下文丢失现象任务函数里写了await session.commit()日志显示Sync completed successfully但数据库表无变化。根本原因FastAPI 的get_async_session是依赖注入返回的 session 对象绑定在当前请求的asyncio.Task上。而 APScheduler 的 job 是在独立的asyncio.Task中执行没有 FastAPI 的依赖注入上下文get_async_session()返回的 session 未正确初始化。解决方案在任务函数中必须显式创建新的 session不能依赖全局注入# ❌ 错误试图用 FastAPI 的依赖注入 async def bad_task(): async for session in get_async_session(): # 这里会报错RuntimeError: async generator not awaited ... # ✅ 正确手动创建 session from app.database import AsyncSessionLocal async def good_task(): async with AsyncSessionLocal() as session: # AsyncSessionLocal 是 sqlalchemy.ext.asyncio.AsyncSession 的实例 # 执行数据库操作 await session.execute(...) await session.commit()AsyncSessionLocal的定义示例# database.py from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession from sqlalchemy.orm import sessionmaker engine create_async_engine( postgresqlasyncpg://user:passlocalhost/db, echoFalse, pool_pre_pingTrue, pool_recycle3600, ) AsyncSessionLocal sessionmaker( engine, class_AsyncSession, expire_on_commitFalse )5.3 “任务重复执行”——RedisJobStore 配置或网络问题现象同一任务在日志中出现两次执行记录时间几乎相同。排查清单可能原因检查方法解决方案Redis 连接不稳定redis-cli ping是否超时redis-cli info clients查看 connected_clients 数量是否突增增加 Redis 连接池大小或检查网络抖动多个 FastapiAdmin 实例未共享 Redis DBredis-cli -n 0 keys *和redis-cli -n 1 keys *对比确保所有实例的db1一致max_instances设置过大查看TaskLog表同一job_id在同一秒内有多条记录将max_instances降为 1或在任务函数开头加分布式锁Cron 表达式触发频率过高*/1 * * * *每分钟触发但任务执行耗时 60 秒改用interval触发器或增加coalesceTrue我们曾遇到一个案例Redis 主从同步延迟 200ms导致两台实例几乎同时读取到 job 的next_run_time各自认为该执行结果双写。最终方案是在TaskLog表加唯一索引(job_id, executed_at)数据库层面拦截重复。5.4 “任务执行报错但日志里看不到堆栈”——事件监听器未正确 await现象APScheduler 日志显示EVENT_JOB_ERROR但你的TaskLog表里error_message为空或只有None。原因scheduler.add_listener()注册的监听器是同步函数而你在监听器里写了await导致协程未被调度。修复如前文scheduler.py所示监听器内不能await必须用asyncio.create_task()包装# ❌ 错误监听器里直接 await def _job_listener(self, event): await self._save_log_to_db(...) # 这行会报错RuntimeError: no running event loop # ✅ 正确用 create_task 启动新协程 def _job_listener(self, event): asyncio.create_task(self._save_log_to_db(...)) # 立即返回不阻塞5.5 生产环境必须做的 5 项加固措施基于 37 个项目上线后的经验列出最易被忽视但至关重要的加固点Redis 连接健康检查在lifespan startup中添加 Redis 连通性测试try: await redis_client.ping() except Exception as e: logging.critical(fRedis connection failed: {e}) raise任务执行超时控制APScheduler 本身不提供超时需在任务函数内手动控制async def task_with_timeout(): try: await asyncio.wait_for(long_running_operation(), timeout300) # 5分钟超时 except asyncio.TimeoutError: logger.error(Task timed out after 300 seconds) raiseRedis JobStore 的 TTL 设置默认 job 数据永不过期长期运行后 Redis 内存暴涨。在RedisJobStore初始化时加jobstores { default: RedisJobStore( ..., # 设置 job 元数据 30 天后自动过期 job_ttl30*24*3600, ) }调度器崩溃自动恢复在lifespan shutdown后加一个守护进程检测调度器状态但更简单的是用 systemd 或 supervisor 管理uvicorn进程崩溃后自动重启APScheduler 会重新加载 Redis 中的任务。Cron 表达式语法校验前置在 FastapiAdmin 后台Task模型的cron_expr字段加 Pydantic validatorfrom croniter import croniter from datetime import datetime field_validator(cron_expr) def validate_cron(cls, v): try: croniter(v, datetime.now()) except Exception as e: raise ValueError(fInvalid cron expression: {e}) return v最后分享一个小技巧我们给每个任务加了tags字段如[data-sync, critical]并在TaskLog表里建了status和tag的复合索引。这样当线上报警“关键任务失败”时DBA 可以 3 秒内查出所有tagcritical且statusfailed的任务比翻日志快 10 倍。这个设计源于一次支付对账任务失败我们花了 47 分钟才定位到是 Redis 连接池耗尽之后就加了这套标签体系。
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表