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

如何将Celery主进程关闭通知传递给ForkPoolWorker实现任务优雅退出?

Celery ForkPoolWorker 优雅退出:主进程关闭通知传递方案

针对主进程关闭通知无法传递给ForkPoolWorker的问题,以下是几个实用的解决思路:

方法一:利用进程共享标志(推荐)

由于ForkPoolWorker是主进程fork生成的,可借助multiprocessing的共享内存对象实现跨进程状态同步:

  1. 初始化共享标志:在Celery应用启动时创建全局布尔共享变量
from multiprocessing import Value
import celery

# 主进程与子进程共享的退出标志
worker_shutting_down = Value('b', False)

app = celery.Celery('your_app')
  1. 主进程处理关闭信号时设置标志
@worker_shutting_down.connect
def handle_worker_shutdown(sig, how, exitcode, **kwargs):
    with worker_shutting_down.get_lock():
        worker_shutting_down.value = True
  1. 任务循环中检查标志并优雅退出
@app.task
def long_running_task():
    for _ in range(100000):
        # 检查共享标志
        with worker_shutting_down.get_lock():
            if worker_shutting_down.value:
                # 执行清理逻辑
                cleanup_resources()
                return

        # 正常执行任务步骤
        execute_task_step()

方法二:主进程转发自定义信号给子进程

主进程捕获关闭信号后,向所有ForkPoolWorker子进程发送自定义信号,子进程捕获后触发退出逻辑:

  1. 主进程转发信号
import os
import signal
from celery.signals import worker_shutting_down

@worker_shutting_down.connect
def forward_shutdown_signal(sig, how, exitcode, **kwargs):
    # 获取所有ForkPoolWorker子进程PID(需根据Celery版本调整筛选逻辑)
    for child_pid in get_worker_child_pids():
        try:
            os.kill(child_pid, signal.SIGUSR1)
        except OSError:
            pass
  1. 子进程注册信号处理函数
import signal
import threading

# 子进程本地退出标志
local_shutdown_flag = threading.Event()

def handle_usr1_signal(signum, frame):
    local_shutdown_flag.set()

# 任务启动前注册信号处理
@task_prerun.connect
def register_signal_handler(sender=None, **kwargs):
    signal.signal(signal.SIGUSR1, handle_usr1_signal)
  1. 任务中检查退出标志
@app.task
def long_running_task():
    for _ in range(100000):
        if local_shutdown_flag.is_set():
            cleanup_resources()
            return

        execute_task_step()

方法三:检查Worker内部状态

任务中定期查询当前worker的运行状态,判断是否处于关闭流程:

from celery import current_task

@app.task
def long_running_task():
    for _ in range(100000):
        worker = current_task.request.worker
        if worker and worker.state in ('shutdown', 'terminated'):
            cleanup_resources()
            return

        execute_task_step()

注意:此方法依赖Celery内部状态的暴露程度,部分版本可能无法直接获取,性能开销略高于前两种方法


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 09:01:09