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

如何在指定Celery任务执行后关闭Worker(避免后续任务)

解决Celery Solo Worker在特定任务完成后关闭的需求

我来帮你搞定这个需求!既然你用的是solo进程池(concurrency=1),我们可以通过Celery的信号机制或者任务内触发关闭逻辑,确保Worker在特定任务执行完毕后不再接收新任务并优雅退出。

核心思路

因为是单进程的solo模式,我们只需要在目标任务完成后触发Worker的关闭信号即可。这里推荐用task_postrun信号(任务执行完成后触发),这种方式能确保任务完全执行完毕再启动关闭流程,比在任务内部直接抛出异常更稳妥。

完整实现代码

from __future__ import absolute_import, unicode_literals
from celery import Celery
from celery.exceptions import WorkerShutdown
from celery.signals import task_postrun

# 初始化Celery实例
app = Celery()
app.config_from_object('celeryconfig')

# 定义你需要触发Worker关闭的特定任务
@app.task(name='target_shutdown_task')
def target_shutdown_task():
    # 这里编写你的任务业务逻辑
    print("特定任务执行完成,即将关闭Worker...")
    return "任务执行成功"

# 绑定任务完成后的信号处理函数
@task_postrun.connect(sender='target_shutdown_task')
def shutdown_worker_after_task(sender, **kwargs):
    """当指定任务执行完成后,触发Worker优雅关闭"""
    # 抛出Celery内置的WorkerShutdown异常,通知Worker停止接收新任务并退出
    raise WorkerShutdown()

关键细节说明

  • 信号精准绑定:通过sender='target_shutdown_task'指定只监听目标任务的完成信号,避免其他任务误触发Worker关闭。
  • WorkerShutdown异常:这是Celery内置的异常类型,抛出后Worker会立刻停止接收新任务,在当前任务(已完成)处理结束后直接退出,完全符合你的需求。
  • solo模式适配:因为是单进程架构,抛出该异常会直接终止Worker进程,不会出现子进程残留的问题。

验证方式

启动Worker时使用以下命令:

celery -A your_app_name worker --pool=solo --concurrency=1 --loglevel=info

当你调用target_shutdown_task.delay()后,任务执行完成,Worker会输出类似Worker shutting down的日志,随后退出,不再接收任何新任务。

备选方案:任务内部直接触发

如果你不需要复用关闭逻辑,也可以直接在任务内部触发关闭(适合简单场景):

@app.task(name='target_shutdown_task')
def target_shutdown_task():
    # 任务业务逻辑
    print("特定任务执行完成,准备关闭Worker...")
    # 直接抛出关闭异常
    raise WorkerShutdown()

注意:如果该任务设置了重试机制,这种方式可能会在重试前就触发关闭,因此信号方式的兼容性更好。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 10:33:04