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

Celery技术问题:如何在指定时间终止运行中的任务?

问题:确保Celery任务在提交后指定时间终止

场景:过载的Celery Worker(设置concurrency=1模拟业务负载),任务无法立即被处理。

模拟繁忙Worker

uv run celery --app src.tasks.example worker --concurrency 1 --queues celery --loglevel info --hostname default 

测试任务定义

@app.task()
def naptime(seconds: float) -> None:
    sleep(seconds)

任务提交方式

naptime.apply_async(args=(3,))
task2 = naptime.apply_async(args=(3,))
start = time()
result = task2.get()
total = time() - start

预期total约为6秒(两个任务顺序执行),实际结果符合预期。

现在需要实现:task2在提交后5秒若仍在运行则被终止

尝试过的无效方案:

  • task2 = naptime.apply_async(args=(3,), time_limit=5):仅在Worker拾取任务后开始计时5秒终止,此时任务可能已完成,但总耗时已超过提交后5秒,不符合需求。
  • task2 = naptime.apply_async(args=(3,), time_limit=5, expires=5):任务来不及过期,无法满足需求。

注:实际场景中无法预知任务执行时长S,仅知道提交后需终止的时间T,因此无法计算expires(即T-S)。


解决方案

核心思路是从任务提交时间开始计时,到指定时间后主动检查并终止任务,可通过以下两种方式实现:

方式1:本地线程监控任务

在提交任务的进程中启动单独线程,等待指定时间后检查任务状态,若任务仍在运行则调用revoke终止:

import time
from threading import Thread
from celery.result import AsyncResult

def terminate_task_after(task_id, delay):
    time.sleep(delay)
    result = AsyncResult(task_id)
    # 检查任务是否待执行或正在执行
    if result.state in ['PENDING', 'STARTED']:
        # terminate=True强制杀死Worker进程中的任务
        app.control.revoke(task_id, terminate=True)

# 提交目标任务
task2 = naptime.apply_async(args=(3,))
# 启动监控线程,5秒后检查并终止任务
Thread(target=terminate_task_after, args=(task2.id, 5), daemon=True).start()

start = time()
try:
    result = task2.get(timeout=5)
except Exception as e:
    print(f"任务超时或被终止: {e}")
total = time() - start

方式2:Celery定时任务触发终止

如果提交任务的进程可能提前退出,可借助Celery定时任务实现延迟终止:

from datetime import datetime, timedelta

# 定义终止任务的函数
@app.task
def terminate_task(task_id):
    result = AsyncResult(task_id)
    if result.state in ['PENDING', 'STARTED']:
        app.control.revoke(task_id, terminate=True)

# 提交目标任务
task2 = naptime.apply_async(args=(3,))
# 计算终止时间:提交时间+5秒
terminate_time = datetime.now() + timedelta(seconds=5)
# 提交定时终止任务,指定执行时间
terminate_task.apply_async(args=(task2.id,), eta=terminate_time)

start = time()
try:
    result = task2.get(timeout=5)
except Exception as e:
    print(f"任务超时或被终止: {e}")
total = time() - start

关键说明

  • app.control.revoke(task_id, terminate=True):必须添加terminate=True才能强制终止Worker中正在运行的任务,仅调用revoke只会阻止未被拾取的任务执行。
  • 两种方案均无需预知任务执行时长,完全基于任务提交时间计算终止节点,匹配需求。
  • 使用定时任务方式时,需确保Celery Beat服务正在运行,启动命令:
    uv run celery --app src.tasks.example beat --loglevel info
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 08:52:39