You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Celery测试Worker无法在超时内停止问题求助

解决Celery Worker定时器导致pytest-celery测试无法正常退出的问题

核心原因

你通过worker_ready信号注册的call_repeatedly定时器会创建持续运行的线程,Celery Worker默认不会自动清理自定义添加的定时器,导致线程一直存活,阻止Worker进程正常退出,最终触发测试超时错误。

解决方案

1. 监听Worker关闭信号,主动取消定时器

绑定worker_shutdown信号回调,手动取消之前注册的定时器条目。call_repeatedly会返回一个TimerEntry对象,保存该对象并在Worker关闭时调用cancel()方法即可:

from celery.signals import worker_ready, worker_shutdown
from celery.worker.consumer import Consumer
from celery.utils.timer import Timer

# 用全局变量存储定时器条目
_timer_entry = None

@worker_ready.connect
def worker_ready_handler(sender: Consumer, **kwargs):
    global _timer_entry
    logger.info(f"Worker process ready: {sender.pid}")
    timer: Timer = sender.timer
    _timer_entry = timer.call_repeatedly(60, do_something, ())

@worker_shutdown.connect
def worker_shutdown_handler(sender: Consumer, **kwargs):
    global _timer_entry
    if _timer_entry:
        _timer_entry.cancel()
        _timer_entry = None

如果不想用全局变量,可将定时器条目绑定到sender实例上,更优雅:

@worker_ready.connect
def worker_ready_handler(sender: Consumer, **kwargs):
    logger.info(f"Worker process ready: {sender.pid}")
    timer: Timer = sender.timer
    # 将定时器条目绑定到sender的自定义属性
    sender._custom_timer = timer.call_repeatedly(60, do_something, ())

@worker_shutdown.connect
def worker_shutdown_handler(sender: Consumer, **kwargs):
    if hasattr(sender, "_custom_timer"):
        sender._custom_timer.cancel()
        del sender._custom_timer

2. 测试环境下直接禁用定时器

如果测试不需要该定时任务功能,可在测试阶段跳过定时器注册:

方式一:环境变量控制

在原代码中添加环境判断:

import os

@worker_ready.connect
def worker_ready_handler(sender: Consumer, **kwargs):
    # 测试环境下不启动定时器
    if os.getenv("TESTING") == "1":
        return
    logger.info(f"Worker process ready: {sender.pid}")
    timer: Timer = sender.timer
    timer.call_repeatedly(60, do_something, ())

在pytest.ini配置文件中设置测试环境变量:

[pytest]
env =
    TESTING=1

方式二:测试中移除信号接收器

在测试fixture中临时移除自定义的worker_ready处理函数:

import pytest
from celery.signals import worker_ready

@pytest.fixture(autouse=True)
def disable_custom_timer():
    # 复制接收器列表,避免遍历过程中修改原列表
    receivers = worker_ready.receivers.copy()
    # 找到并移除目标处理函数
    for receiver in receivers:
        if receiver.__name__ == "worker_ready_handler":
            worker_ready.disconnect(receiver)
    yield
    # 测试结束后恢复(可选,因每个测试用例的Worker是独立实例)
    worker_ready.connect(worker_ready_handler)

3. 确保定时任务函数可中断

如果do_something函数包含阻塞逻辑(如长时间循环、IO操作),即使取消定时器,正在运行的函数实例也可能无法终止。需让函数能响应退出信号:

import threading
import time

def do_something():
    while threading.current_thread().is_alive():
        # 执行你的任务逻辑
        # ...
        
        # 添加短睡眠,让线程有机会响应取消操作
        time.sleep(1)
        
        # 可选:通过线程属性设置退出标志
        if getattr(threading.current_thread(), "should_exit", False):
            break

验证

修改后运行多个测试用例,Worker应能在测试结束后正常退出,不会再触发超时错误。

内容的提问来源于stack exchange,提问作者tenticon

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 18:35:26