配置开发环境 Celery 在 FastAPI 的 lifespan 启动时自动开启,结束时自动关闭;生产环境中 Celery Worker 的独立进程启动方式
在 FastAPI 的lifespan中直接启动 Celery Worker 仅适用于开发环境或单进程演示场景。生产环境中,Celery Worker 必须作为独立进程运行,否则会导致资源竞争、信号处理异常及性能瓶颈。
创建celeryconfig.py文件
import logging import threading from celery import Celery from celery.worker import state as worker_state from config.configure import ( CELERY_CONCURRENCY, REDIS_DB, REDIS_HOST, REDIS_PORT, ) logger = logging.getLogger(__name__) REDIS_URL = f"redis://{REDIS_HOST}:{REDIS_PORT}" celery_app = Celery( "python_mysql", broker=f"{REDIS_URL}/{REDIS_DB}", backend=f"{REDIS_URL}/{REDIS_DB + 1}", # 任务结果存相邻 DB,避免与业务缓存混库 include=[__name__], # worker 启动时自动导入本模块注册任务( # __name__ 同时兼容 celery -A celeryconfig # 与 config.celeryconfig 两种导入方式) ) celery_app.conf.update( timezone="Asia/Shanghai", enable_utc=False, task_track_started=True, task_acks_late=True, # 任务执行成功后才确认,worker 崩溃时任务重新投递 worker_prefetch_multiplier=1, # 配合 acks_late,避免任务预取堆积 broker_connection_retry_on_startup=True, **({"worker_concurrency": CELERY_CONCURRENCY} if CELERY_CONCURRENCY else {}), ) @celery_app.task( autoretry_for=(Exception,), # 抛出指定异常时自动重试 retry_backoff=5, # 指数退避:第 n 次重试等待 5 * 2^(n-1) 秒 retry_backoff_max=60, # 退避上限 60 秒 retry_jitter=True, # 加入随机抖动,避免重试雪崩 retry_kwargs={"max_retries": 3}, # 最多重试 3 次 ) def example_task(x: int, y: int) -> int: """示例任务:两数相加。除 0 触发异常可观察自动重试行为""" if x == 0 and y == 0: raise ValueError("x 和 y 不能同时为 0") # 用于演示自动重试 result = x + y logger.info("example_task: %s + %s = %s", x, y, result) return result # ---------------------------------------------------------------------------- # 内嵌 worker:随 FastAPI lifespan 自动启动/关闭,无需单独执行 celery 命令 # ---------------------------------------------------------------------------- # macOS 上 billiard 默认 spawn 子进程,内嵌时用 solo 池(同进程线程内执行)最稳妥 EMBEDDED_WORKER_POOL = "solo" EMBEDDED_WORKER_CONCURRENCY = 1 _embedded_worker = None _worker_thread = None _worker_ready = threading.Event() def _run_embedded_worker() -> None: """在独立守护线程中构造并运行 worker,阻塞至 worker 退出。""" global _embedded_worker try: worker = celery_app.Worker( loglevel="info", pool=EMBEDDED_WORKER_POOL, concurrency=EMBEDDED_WORKER_CONCURRENCY, ) _embedded_worker = worker _worker_ready.set() # 通知主线程:worker 已构造完成 worker.start() # 阻塞运行,stop() 后返回 except Exception: logger.exception("内嵌 Celery worker 运行异常") _worker_ready.set() def start_embedded_worker(timeout: float = 10.0) -> None: """FastAPI lifespan 启动时调用:在守护线程中拉起 Celery worker。""" global _worker_thread if _worker_thread is not None and _worker_thread.is_alive(): logger.warning("Celery worker 已在运行,跳过重复启动") return worker_state.should_stop = None # 复位关闭标志,防止历史状态残留 _worker_ready.clear() _worker_thread = threading.Thread( target=_run_embedded_worker, name="celery-worker", daemon=True, # 守护线程:异常情况下也不会拖住进程退出 ) _worker_thread.start() _worker_ready.wait(timeout=timeout) if _embedded_worker is None: logger.error("Celery 内嵌 worker 启动失败,请检查上方异常日志") else: logger.info("Celery 内嵌 worker 已启动(pool=%s)", EMBEDDED_WORKER_POOL) def stop_embedded_worker(timeout: float = 10.0) -> None: """FastAPI lifespan 关闭时调用:优雅停止 worker 并等待线程退出。""" global _embedded_worker, _worker_thread worker, thread = _embedded_worker, _worker_thread if worker is None or thread is None: return # 置位 should_stop:事件循环与 broker 重连循环会据此退出, # 即便 worker 尚处在启动(连接 broker)阶段也能被打断 worker_state.should_stop = True try: worker.stop() # 优雅关闭(warm shutdown) except Exception: logger.exception("优雅关闭 Celery worker 失败,尝试强制终止") worker.terminate() thread.join(timeout=timeout) if thread.is_alive(): logger.warning("Celery worker 未在 %ss 内退出,守护线程将随进程结束", timeout) _embedded_worker = None _worker_thread = None logger.info("Celery 内嵌 worker 已关闭") """ 方式一(推荐,已接入 main.py 的 lifespan): 随 FastAPI 启动自动开启、关闭自动结束,无需手动管理进程。 方式二(独立进程,项目根目录执行): celery -A config.celeryconfig worker -l info -P solo 注意:macOS 上 billiard 默认 spawn 子进程,prefork 池无法继承任务注册表, 必须用 -P solo(单进程,本地验证)或 -P threads(多线程,支持并发)。 调用任务(Python 侧): from config.celeryconfig import example_task result = example_task.delay(1, 2) # 异步投递 print(result.get(timeout=10)) # 阻塞获取结果 """在main.py里
from contextlib import asynccontextmanager from config.database import init_db, close_db from config.redis import close_redis @asynccontextmanager async def life_span(app: FastAPI): print("server is starting …") await init_db() # 仅开发模式(CELERY_EMBEDDED=true,默认)内嵌 worker; # 生产环境使用独立进程:scripts/celery_worker.sh if CELERY_EMBEDDED: start_embedded_worker() else: logging.info( "CELERY_EMBEDDED=false,跳过内嵌 worker,请确认独立 Celery 进程已启动" ) yield # 先停 Celery(需使用 Redis broker),再关闭 Redis/DB 连接 if CELERY_EMBEDDED: stop_embedded_worker() await close_db() await close_redis() print("server has been shut down.") app = FastAPI(lifespan=life_span)在根目录创建脚本文件celery_worker.sh和.env
celery_worker.sh
#!/usr/bin/env bash # ============================================================================= # Celery worker 独立进程管理脚本(生产部署 / macOS 本地验证均可用) # # 用法(项目根目录执行): # scripts/celery_worker.sh start # 后台启动 worker(单实例,重复启动会被拒绝) # scripts/celery_worker.sh stop # 优雅停止(SIGTERM warm shutdown,超时再 KILL) # scripts/celery_worker.sh restart # 等同于 stop + start # scripts/celery_worker.sh status # 查看运行状态(退出码 0=运行中,3=未运行) # # 可选环境变量(也可写在 .env 中由 config.configure 读取,CLI 参数优先级更高): # CELERY_POOL=solo|prefork|threads # 不设置时:Linux=prefork(默认),macOS=solo # CELERY_CONCURRENCY=4 # worker 并发数,默认按 CPU 核数 # CELERY_QUEUES=celery # 监听队列,多个用逗号分隔 # CELERY_HOSTNAME=worker@%h-1 # 节点名,多实例部署时区分 # CELERY_LOG_LEVEL=info # 日志级别 # # 退出码:0 成功;1 参数或启动失败;3 未在运行(status) # ============================================================================= set -euo pipefail APP="config.celeryconfig" # 定位项目根目录(脚本位于 <root>/scripts/ 下) PROJECT_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" cd "$PROJECT_ROOT" # 优先使用项目虚拟环境内的 celery,其次使用 PATH 中的 CELERY_BIN="$PROJECT_ROOT/venv/bin/celery" if [ ! -x "$CELERY_BIN" ]; then CELERY_BIN="$(command -v celery)" fi if [ -z "$CELERY_BIN" ]; then echo "未找到 celery 可执行文件,请先安装依赖或激活虚拟环境" >&2 exit 1 fi # 运行参数 -------------------------------------------------------------- POOL="${CELERY_POOL:-}" # macOS 上 billiard 默认 spawn,prefork 无法继承任务注册表,自动回退 solo if [ -z "$POOL" ] && [ "$(uname -s)" = "Darwin" ]; then POOL="solo" fi CONCURRENCY="${CELERY_CONCURRENCY:-}" QUEUES="${CELERY_QUEUES:-celery}" HOSTNAME="${CELERY_HOSTNAME:-}" LOG_LEVEL="${CELERY_LOG_LEVEL:-info}" RUN_DIR="$PROJECT_ROOT/run" LOG_DIR="$PROJECT_ROOT/logs" PIDFILE="$RUN_DIR/celery_worker.pid" LOGFILE="$LOG_DIR/celery_worker.log" # 停止时最长等待秒数(warm shutdown,prefork 下会等在途任务执行完) STOP_TIMEOUT=30 mkdir -p "$RUN_DIR" "$LOG_DIR" is_running() { [ -f "$PIDFILE" ] || return 1 local pid pid="$(cat "$PIDFILE" 2>/dev/null || true)" [ -n "$pid" ] && kill -0 "$pid" 2>/dev/null } build_args() { ARGS=("-A" "$APP" "worker" "-l" "$LOG_LEVEL" "-Q" "$QUEUES" "--pidfile" "$PIDFILE" "--logfile" "$LOGFILE") if [ -n "$POOL" ]; then ARGS+=("-P" "$POOL") fi if [ -n "$CONCURRENCY" ]; then ARGS+=("-c" "$CONCURRENCY") fi if [ -n "$HOSTNAME" ]; then ARGS+=("-n" "$HOSTNAME") fi } start_worker() { if is_running; then echo "celery worker 已在运行 (pid=$(cat "$PIDFILE")),跳过启动" return 0 fi # 清理上次异常退出残留的失效 pidfile rm -f "$PIDFILE" build_args echo "启动 celery worker: $CELERY_BIN ${ARGS[*]}" # nohup 后台运行,stdout/stderr 追加到同一日志文件(承接 celery logfile 之外的早期输出) nohup "$CELERY_BIN" "${ARGS[@]}" >> "$LOGFILE" 2>&1 & # 等待 celery 写入 pidfile 且进程存活(约 10s) for _ in $(seq 1 50); do if is_running; then echo "celery worker 已启动 (pid=$(cat "$PIDFILE")),日志: $LOGFILE" return 0 fi sleep 0.2 done echo "celery worker 启动失败,请查看日志: $LOGFILE" >&2 exit 1 } stop_worker() { if ! is_running; then echo "celery worker 未在运行" rm -f "$PIDFILE" return 0 fi local pid pid="$(cat "$PIDFILE")" echo "优雅停止 celery worker (pid=$pid),最长等待 ${STOP_TIMEOUT}s ..." # SIGTERM:Celery warm shutdown,等待在途任务执行完毕后退出 kill -TERM "$pid" 2>/dev/null || true local waited=0 while kill -0 "$pid" 2>/dev/null; do if [ "$waited" -ge "$STOP_TIMEOUT" ]; then echo "等待超时,强制结束 worker 进程组 ..." >&2 # 兜底:仅精确匹配本项目虚拟环境 + 本 app 的 worker 进程,避免误杀其他项目 pkill -KILL -f "$CELERY_BIN.*-A $APP worker" 2>/dev/null || true break fi sleep 1 waited=$((waited + 1)) done rm -f "$PIDFILE" echo "celery worker 已停止" } status_worker() { if is_running; then echo "celery worker 运行中 (pid=$(cat "$PIDFILE"))" exit 0 fi echo "celery worker 未运行" exit 3 } case "${1:-}" in start) start_worker ;; stop) stop_worker ;; restart) stop_worker start_worker ;; status) status_worker ;; *) echo "用法: $0 {start|stop|restart|status}" >&2 exit 1 ;; esac.env
REDIS_HOST=localhost REDIS_PORT=6379 # Celery 运行模式:true=内嵌 worker(本地开发,随 FastAPI 自动启停) # 生产环境改为 false,worker 用 scripts/celery_worker.sh 独立进程管理 CELERY_EMBEDDED=true