任务
任务是 Celery 应用的基本构成单元。
任务是一个类,可以从任何可调用对象创建。它承担两种职责:定义调用任务时发生什么(发送消息),以及 worker 收到该消息后发生什么。
每个任务类都有唯一名称。消息引用这个名称,worker 才能找到正确的函数执行。
任务消息只有在被 worker 确认后才会从队列删除。worker 可以预先保留多条消息;即使它因断电或其他原因被终止,消息也会重新投递给其他 worker。
理想情况下,任务函数应当幂等:即使使用相同参数调用多次,也不会产生非预期影响。worker 无法判断任务是否幂等,因此默认在执行前确认消息,让已经开始执行的任务调用不再执行第二次。
如果任务具有幂等性,可以设置 acks_late,改为在任务返回后确认消息。另请参阅常见问题应该使用retry还是acks_late?。
即使启用了 acks_late,执行任务的子进程终止时,无论由任务调用 sys.exit() 还是由信号触发,worker 仍会确认消息。这是有意设计的,因为:
-
不希望重新运行会让内核向进程发送
SIGSEGV(段错误)或类似信号的任务。 -
系统管理员主动终止任务时,通常并不希望它自动重启。
-
分配过多内存的任务可能触发内核 OOM killer;重新运行时,同样的问题可能再次发生。
-
每次重新投递都失败的任务,可能形成高频消息循环并拖垮系统。
如果确实希望在这些场景中重新投递任务,可以考虑启用 task_reject_on_worker_lost。
警告
无限期阻塞的任务最终可能使整个 worker 实例无法再处理任何其他工作。
任务涉及I/O时,务必为这些操作设置超时。例如,用 https://pypi.org/project/requests/ 库发送网页请求时添加超时:
connect_timeout, read_timeout = 5.0, 30.0
response = requests.get(URL, timeout=(connect_timeout, read_timeout))
时间限制有助于确保所有任务及时返回,但超时事件实际上会强制杀死进程,因此应仅用它们检测尚未手动设置超时的情况。
旧版本中,默认 prefork 进程池调度器不适合长时间运行的任务,因此对于运行数分钟或数小时的任务,曾建议在 celery worker 上启用 -Ofair 命令行参数。从4.0开始,-Ofair 已成为默认调度策略。更多信息见预取限制。为获得最佳性能,应将长任务和短任务路由到各自专用的 worker,参见自动路由。
worker 卡住时,请在提交问题之前调查哪些任务正在执行,因为最常见原因是一个或多个任务卡在网络操作上。
—
本章介绍任务定义的各个方面,目录如下:
基础
使用 app.task() 装饰器,可以方便地从任意可调用对象创建任务:
from .models import User
@app.task
def create_user(username, password):
User.objects.create(username=username, password=password)
任务还有许多可配置的选项,可以作为参数传给装饰器:
@app.task(serializer='json')
def create_user(username, password):
User.objects.create(username=username, password=password)
如何导入任务装饰器?
多个装饰器
将任务装饰器与其他装饰器组合使用时,必须确保任务装饰器最后应用。Python 的写法看起来相反:它必须放在装饰器列表的最上方。
@app.task @decorator2 @decorator1 def add(x, y): return x + y
绑定任务
绑定任务意味着第一个参数始终是任务实例 self,类似于 Python 绑定方法:
logger = get_task_logger(__name__)
@app.task(bind=True)
def add(self, x, y):
logger.info(self.request.id)
重试(使用 app.Task.retry())、访问当前任务请求的信息,以及使用自定义任务基类添加的其他功能时,都需要绑定任务。
任务继承
任务装饰器的 base 参数指定任务基类:
import celery
class MyTask(celery.Task):
def on_failure(self, exc, task_id, args, kwargs, einfo):
print('{0!r} failed: {1!r}'.format(task_id, exc))
@app.task(base=MyTask)
def add(x, y):
raise KeyError()
名称
每个任务都必须有唯一名称。
没有显式指定名称时,任务装饰器会自动生成名称,依据是任务定义所在的模块,以及任务函数名称。
显式指定名称的示例:
>>> @app.task(name='sum-of-two-numbers')
>>> def add(x, y):
... return x + y
>>> add.name
'sum-of-two-numbers'
建议将模块名用作命名空间。即使其他模块已经定义同名任务,也不会发生名称冲突。
>>> @app.task(name='tasks.add')
>>> def add(x, y):
... return x + y
可以通过任务的 .name 属性查看名称:
>>> add.name
'tasks.add'
这里指定的名称 tasks.add,恰好就是在名为 tasks.py 的模块中定义此任务时自动生成的名称:
tasks.py:
@app.task
def add(x, y):
return x + y
>>> from tasks import add
>>> add.name
'tasks.add'
注意
可以对 worker 使用 inspect 命令,查看所有已注册任务名称。参见用户指南管理命令行工具(inspect/control)中的 inspect registered 命令。
修改自动命名行为
4.0版本新增。
默认自动命名并不总是合适。例如,多个模块中定义了许多任务:
project/
/__init__.py
/celery.py
/moduleA/
/__init__.py
/tasks.py
/moduleB/
/__init__.py
/tasks.py
默认情况下,任务名称会类似 moduleA.tasks.taskA、moduleA.tasks.taskB、moduleB.tasks.test。可能希望从所有任务名称中去掉 tasks。除了为每个任务显式命名,也可以覆写 app.gen_task_name(),改变自动命名行为。继续上面的例子,celery.py 可以包含:
from celery import Celery
class MyCelery(Celery):
def gen_task_name(self, name, module):
if module.endswith('.tasks'):
module = module[:-6]
return super().gen_task_name(name, module)
app = MyCelery('main')
各任务名称就会变为 moduleA.taskA、moduleA.taskB、moduleB.test。
警告
确保 app.gen_task_name() 是纯函数:相同输入必须始终返回相同输出。
任务请求
app.Task.request 包含当前正在执行的任务的相关信息和状态。
请求定义以下属性:
- id:
-
正在执行的任务的唯一ID。
- group:
-
如果任务属于一个group,这里是该组的唯一ID。
- chord:
-
任务所属 chord 的唯一ID,适用于任务属于 header 的情况。
- correlation_id:
-
可用于去重等用途的自定义ID。
- args:
-
位置参数。
- kwargs:
-
关键字参数。
- origin:
-
发送该任务的主机名称。
- retries:
-
当前任务已经重试的次数,是从0开始的整数。
- is_eager:
-
如果任务在客户端本地执行,而不是由 worker 执行,此值为
True。 - eta:
-
任务最初的预计执行时间ETA(如果有)。使用UTC时间,具体取决于
enable_utc设置。 - expires:
-
任务最初的过期时间(如果有)。使用UTC时间,具体取决于
enable_utc设置。 - hostname:
-
执行任务的 worker 实例的节点名称。
- delivery_info:
-
额外的消息投递信息,是包含投递此任务所用交换机及路由键的映射。例如
app.Task.retry()会用它将任务重新发送到相同的目标队列。该字典中有哪些键取决于使用的消息代理。 - reply-to:
-
接收回复的队列名称,例如RPC结果后端会使用它。
- called_directly:
-
任务不是由 worker 执行时,此标志为true。
- timelimit:
-
当前对该任务生效的
(soft, hard)时间限制元组(如果有)。 - callbacks:
-
任务成功返回时要调用的签名列表。
- errbacks:
-
任务失败时要调用的签名列表。
- utc:
-
调用方启用UTC时为true,参见
enable_utc。
3.1版本新增。
- headers:
-
随任务消息发送的消息头映射,可以为
None。 - reply_to:
-
发送回复的目的地,即队列名称。
- correlation_id:
-
通常与任务ID相同;AMQP中常用它跟踪回复所对应的请求。
4.0版本新增。
- root_id:
-
该任务所属工作流中第一个任务的唯一ID(如果有)。
- parent_id:
-
调用当前任务的父任务的唯一ID(如果有)。
- chain:
-
组成链的任务的逆序列表(如果有)。列表最后一项是当前任务之后要执行的任务。使用第一版任务协议时,链中的任务改为保存在
request.callbacks中。
5.2版本新增。
- properties:
-
随任务消息接收的消息属性映射,可以为
None或{}。 - replaced_task_nesting:
-
任务被替换的次数(如果发生过),可以为
0。
示例
下面的任务访问上下文中的信息:
@app.task(bind=True)
def dump_context(self, x, y):
print('Executing task id {0.id}, args: {0.args!r} kwargs: {0.kwargs!r}'.format(
self.request))
bind 参数让函数成为“绑定方法”,从而可以访问任务类型实例的属性和方法。
日志
worker 会自动设置日志,也可以手动配置。
名为 celery.task 的专用日志记录器可供使用。继承它可以在日志中自动包含任务名称和唯一ID。
建议在模块顶部为所有任务创建一个共用日志记录器:
from celery.utils.log import get_task_logger
logger = get_task_logger(__name__)
@app.task
def add(x, y):
logger.info('Adding {0} + {1}'.format(x, y))
return x + y
Celery 使用 Python 标准日志库,文档见此处。
也可以使用 print(),因为写入标准输出或标准错误的内容会被重定向到日志系统。可以禁用这一行为,参见worker_redirect_stdouts。
注意
如果在任务或任务模块中的其他位置创建日志记录器实例,worker 不会更新重定向。
要将 sys.stdout 和 sys.stderr 重定向到自定义日志记录器,需要手动启用,例如:
import sys
logger = get_task_logger(__name__)
@app.task(bind=True)
def add(self, x, y):
old_outs = sys.stdout, sys.stderr
rlevel = self.app.conf.worker_redirect_stdouts_level
try:
self.app.log.redirect_stdouts_to_logger(logger, rlevel)
print('Adding {0} + {1}'.format(x, y))
return x + y
finally:
sys.stdout, sys.stderr = old_outs
注意
所需的某个 Celery 日志记录器没有输出日志时,应检查日志是否正确传播。下面启用 celery.app.trace,以输出“succeeded in”日志:
import celery
import logging
@celery.signals.after_setup_logger.connect
def on_after_setup_logger(**kwargs):
logger = logging.getLogger('celery')
logger.propagate = True
logger = logging.getLogger('celery.app.trace')
logger.propagate = True
注意
要完全禁用 Celery 的日志配置,请使用 setup_logging 信号:
import celery
@celery.signals.setup_logging.connect
def on_setup_logging(**kwargs):
pass
参数检查
4.0版本新增。
调用任务时,Celery 会验证传入的参数,类似于 Python 调用普通函数时的行为:
>>> @app.task
... def add(x, y):
... return x + y
# Calling the task with two arguments works:
>>> add.delay(8, 8)
<AsyncResult: f59d71ca-1549-43e0-be41-4e8821a83c0c>
# Calling the task with only one argument fails:
>>> add.delay(8)
Traceback (most recent call last):
File "<stdin>", line 1, in <module>
File "celery/app/task.py", line 376, in delay
return self.apply_async(args, kwargs)
File "celery/app/task.py", line 485, in apply_async
check_arguments(*(args or ()), **(kwargs or {}))
TypeError: add() takes exactly 2 arguments (1 given)
将任务的 typing 属性设为 False,可以关闭参数检查:
>>> @app.task(typing=False)
... def add(x, y):
... return x + y
# Works locally, but the worker receiving the task will raise an error.
>>> add.delay(8)
<AsyncResult: f59d71ca-1549-43e0-be41-4e8821a83c0c>
隐藏参数中的敏感信息
4.0版本新增。
使用第2版或更新的 task_protocol 时(4.0起默认为第2版),可以通过调用参数 argsrepr 和 kwargsrepr,改变位置参数与关键字参数在日志及监控事件中的显示方式:
>>> add.apply_async((2, 3), argsrepr='(<secret-x>, <secret-y>)')
>>> charge.s(account, card='1234 5678 1234 5678').set(
... kwargsrepr=repr({'card': '**** **** **** 5678'})
... ).delay()
警告
能够从消息代理读取任务消息,或者截获消息的人,仍然能够获得其中的敏感信息。
因此,包含敏感信息的消息通常应加密。例如,信用卡号码可以加密保存在安全存储中,由任务本身检索并解密实际号码。
重试
app.Task.retry() 可用于重新执行任务,例如遇到可恢复错误时。
调用 retry 会发送一条新消息,使用相同任务ID,并确保消息投递到原任务所在的同一个队列。
重试也会记录为任务状态,因此可以通过结果实例跟踪任务进度,参见状态。
使用 retry 的示例:
@app.task(bind=True)
def send_twitter_status(self, oauth, tweet):
try:
twitter = Twitter(oauth)
twitter.update_status(tweet)
except (Twitter.FailWhaleError, Twitter.LoginError) as exc:
raise self.retry(exc=exc)
注意
调用 app.Task.retry() 会抛出异常,因此重试调用之后的代码不会执行。抛出的是 Retry,它不被当作错误,而是作为半谓词通知 worker:任务需要重试。在启用结果后端时,worker 因而能保存正确状态。
这是正常行为,除非将 retry 的 throw 参数设为 False,否则总会如此。
任务装饰器的 bind 参数允许访问任务类型实例 self。
exc 参数传递用于日志和任务结果存储的异常信息。启用结果后端时,异常与回溯都可以从任务状态中获取。
如果任务设置了 max_retries,超出最大重试次数时会重新抛出当前异常,但以下情况除外:
-
没有提供
exc参数。此时抛出
MaxRetriesExceededError。 -
当前没有异常。
没有原始异常可以重新抛出时,会使用
exc参数。因此:self.retry(exc=Twitter.LoginError())会抛出传入的
exc。
自定义重试延迟
任务可以等待指定时间后再重试,默认延迟由 default_retry_delay 属性定义,默认值为3分钟。注意,设置延迟时的单位是秒,可以为整数或浮点数。
也可以向 retry() 传入 countdown 参数,覆盖默认值。
@app.task(bind=True, default_retry_delay=30 * 60) # retry in 30 minutes.
def add(self, x, y):
try:
something_raising()
except Exception as exc:
# overrides the default delay to retry after 1 minute
raise self.retry(exc=exc, countdown=60)
针对已知异常自动重试
4.0版本新增。
有时希望只要出现特定异常就重试任务。
可以在 app.task() 装饰器中设置 autoretry_for 参数,让 Celery 自动重试:
from twitter.exceptions import FailWhaleError
@app.task(autoretry_for=(FailWhaleError,))
def refresh_timeline(user):
return twitter.refresh_timeline(user)
若要为内部的 retry() 调用指定自定义参数,可向 app.task() 装饰器传入 retry_kwargs:
@app.task(autoretry_for=(FailWhaleError,),
retry_kwargs={'max_retries': 5})
def refresh_timeline(user):
return twitter.refresh_timeline(user)
这是手动处理异常之外的另一种方式。上例等同于将任务函数体包裹在 try…except 语句中:
@app.task
def refresh_timeline(user):
try:
twitter.refresh_timeline(user)
except FailWhaleError as exc:
raise refresh_timeline.retry(exc=exc, max_retries=5)
若要遇到任何错误都自动重试,只需:
@app.task(autoretry_for=(Exception,))
def x():
...
4.2版本新增。
任务依赖其他服务,例如调用API时,适合使用指数退避,避免大量请求压垮服务。Celery 的自动重试支持让这很容易实现,只需指定 retry_backoff 参数:
from requests.exceptions import RequestException
@app.task(autoretry_for=(RequestException,), retry_backoff=True)
def x():
...
默认情况下,指数退避还会引入随机抖动,避免所有任务同时运行,并将最大退避延迟限制为10分钟。这些设置都可以通过下文选项自定义。
4.4版本新增。
基于类的任务也可以设置 autoretry_for、max_retries、retry_backoff、retry_backoff_max 和 retry_jitter:
class BaseTaskWithRetry(Task):
autoretry_for = (TypeError,)
max_retries = 5
retry_backoff = True
retry_backoff_max = 700
retry_jitter = False
- Task.autoretry_for
-
异常类的列表或元组。任务执行时出现其中任意异常,都会自动重试。默认不会对任何异常自动重试。
- Task.max_retries
-
放弃之前允许的最大重试次数。设为
None表示永远重试,默认值为3。
- Task.retry_backoff
-
布尔值或数字。设为
True时,自动重试按指数退避延迟:第一次等待1秒,第二次2秒,第三次4秒,第四次8秒,以此类推;启用retry_jitter时,实际延迟还会被调整。设为数字时,该数字用作延迟系数。例如设为3时,各次延迟为3、6、12、24秒等。默认设为False,自动重试不附加延迟。
- Task.retry_backoff_max
-
数字。启用
retry_backoff时,指定自动重试之间的最大延迟,单位为秒。默认是600,即10分钟。
- Task.retry_jitter
-
布尔值。抖动 为指数退避引入随机性,避免队列中的任务同时执行。设为
True时,retry_backoff计算的延迟被视为上限,实际延迟是在0到上限之间随机选择的数值。默认是True。
5.3.0版本新增。
- Task.dont_autoretry_for
-
- 异常类的列表或元组;这些异常不会触发自动重试。
-
可以用它排除虽然匹配
autoretry_for、但不希望重试的异常。
使用Pydantic验证参数
5.5.0版本新增。
传入 pydantic=True,可以使用 Pydantic 根据类型提示验证并转换参数,以及序列化结果。
注意
参数验证仅覆盖任务端的参数和返回值。通过 delay() 或 apply_async() 调用任务时,仍需自行序列化参数。
例如:
from pydantic import BaseModel
class ArgModel(BaseModel):
value: int
class ReturnModel(BaseModel):
value: str
@app.task(pydantic=True)
def x(arg: ArgModel) -> ReturnModel:
# args/kwargs type hinted as Pydantic model will be converted
assert isinstance(arg, ArgModel)
# The returned model will be converted to a dict automatically
return ReturnModel(value=f"example: {arg.value}")
之后可以传入符合模型定义的字典来调用任务,返回的模型会被“导出”,即通过 BaseModel.model_dump() 序列化:
>>> result = x.delay({'value': 1})
>>> result.get(timeout=1)
{'value': 'example: 1'}
联合类型与泛型参数
不支持联合类型,例如 Union[SomeModel, OtherModel],也不支持泛型参数,例如 list[SomeModel]。
如果需要支持列表或类似类型,建议使用 pydantic.RootModel。
可选参数与返回值
可选参数及返回值也会正确处理。例如,对于这个任务:
from typing import Optional
# models are the same as above
@app.task(pydantic=True)
def x(arg: Optional[ArgModel] = None) -> Optional[ReturnModel]:
if arg is None:
return None
return ReturnModel(value=f"example: {arg.value}")
会得到以下行为:
>>> result = x.delay()
>>> result.get(timeout=1) is None
True
>>> result = x.delay({'value': 1})
>>> result.get(timeout=1)
{'value': 'example: 1'}
返回值处理
仅当返回的模型与类型注解匹配时,才会序列化返回值。返回其他类型的模型实例时不会序列化。mypy 理应已经能发现这类错误,应修正类型提示。
选项列表
任务装饰器接受多种改变任务行为的选项,例如通过 rate_limit 设置任务速率限制。
传给任务装饰器的任何关键字参数,都会被设置为所生成任务类的属性。下面列出内置属性。
通用选项
- Task.name
-
任务注册时使用的名称。
可以手动设置名称,也可以根据模块和类名自动生成。
另请参阅名称。
- Task.request
-
任务正在执行时,这里包含当前请求的信息,使用线程局部存储。
参见任务请求。
- Task.max_retries
-
只有任务调用
self.retry,或者任务装饰器使用了 autoretry_for 参数时才生效。放弃之前允许尝试的最大重试次数。超过此值时,抛出
MaxRetriesExceededError。注意
需要手动调用
retry(),它不会仅因发生异常就自动重试。默认值为
3。设为None将取消重试次数限制,任务会一直重试直到成功。
- Task.throws
-
可选的预期异常类元组,这些异常不应被视为真正的错误。
这些异常仍会作为失败报告给结果后端,但 worker 不会以错误级别记录该事件,也不会附带回溯。
例如:
@task(throws=(KeyError, HttpNotFound)): def get_foo(): something()错误类型:
-
预期错误,包含在
Task.throws中。以
INFO级别记录,不包含回溯。 -
非预期错误。
以
ERROR级别记录,并包含回溯。
-
- Task.rate_limit
-
设置此任务类型的速率限制,即限定时间范围内能够运行的任务数量。限制生效时,任务仍会执行完成,但开始执行前可能需要等待。
设为
None时不限制速率。整数或浮点数表示每秒任务数。可以在数值后加 /s、/m 或 /h,分别指定每秒、每分钟或每小时的限制。任务会均匀分布到指定时间范围内。
例如100/m表示每分钟100个任务,要求同一 worker 实例启动两个任务之间至少间隔600毫秒。
默认值来自
task_default_rate_limit设置;若未指定,任务默认不限速。注意,这是每个 worker 实例的速率限制,并非全局限制。要实现全局限速,例如API每秒最大请求数,必须将任务限定到指定队列。
- Task.time_limit
-
此任务的硬时间限制,单位为秒。未设置时使用 worker 默认值。
- Task.soft_time_limit
-
此任务的软时间限制。未设置时使用 worker 默认值。
- Task.ignore_result
-
不存储任务状态。这意味着无法通过
AsyncResult检查任务是否完成,也无法获取返回值。注意:禁用任务结果后,某些功能将无法工作,详见 Canvas 文档。
- Task.store_errors_even_if_ignored
-
设为
True时,即使任务被配置为忽略结果,仍会保存错误。5.7版本变更:以前,请求消息缺少
ignore_result键时,store_errors会默认采用True,忽略任务自身的ignore_result设置。现在没有请求级覆盖设置时,worker 会正确回退到Task.ignore_result。
- Task.serializer
-
指定默认序列化方式的字符串,默认采用
task_serializer设置。可以是 pickle、json、yaml,或者通过kombu.serialization.registry注册的自定义序列化方式。更多信息见序列化器。
- Task.compression
-
指定默认压缩方式的字符串。
默认采用
task_compression设置。可以是 gzip、bzip2,或者通过kombu.compression注册表注册的自定义压缩方式。更多信息见压缩。
- Task.backend
-
该任务使用的结果存储后端,是 celery.backends 中某个后端类的实例。默认使用 app.backend,由
result_backend设置定义。
- Task.acks_late
-
设为
True时,任务消息会在执行之后确认,而不是默认的执行之前确认。注意:worker 在执行过程中崩溃时,这意味着任务可能执行多次。请确保任务幂等。
可以通过
task_acks_late设置覆盖全局默认值。
- Task.track_started
-
设为
True时,worker 开始执行任务会报告“started”状态。默认是False,因为通常无需报告这种粒度的状态:任务只需处于等待、完成或等待重试状态。对于长时间运行的任务,如果需要报告哪个任务正在执行,“started”状态会很有用。执行任务的 worker 的主机名及进程ID会包含在状态元数据中,例如 result.info[‘pid’]。
可以通过
task_track_started设置覆盖全局默认值。
另请参阅
Task的API参考。
状态
Celery 可以跟踪任务的当前状态。状态还包含成功任务的结果,或者失败任务的异常及回溯信息。
可以选择不同结果后端,各有优缺点,参见结果后端。
任务在生命周期内可能经历多个状态,各状态都可以附带任意元数据。进入新状态后,旧状态会被遗忘,但某些状态转换仍可推断。例如,当前处于 FAILED 的任务,意味着之前曾处于 STARTED。
另外还有状态集合,例如FAILURE_STATES与READY_STATES。
客户端根据状态是否属于这些集合,决定是否重新抛出异常(PROPAGATE_STATES),以及状态是否可缓存;任务完成后状态可缓存。
也可以定义自定义状态。
结果后端
要跟踪任务或获取返回值,Celery 必须将状态存储或发送到某处,以便稍后检索。内置结果后端包括 SQLAlchemy/Django ORM、Memcached、RabbitMQ/QPid(rpc)及 Redis,也可以自定义后端。
没有任何后端适合所有场景。应了解各后端的优缺点,并选择最符合需求的一种。
警告
后端需要消耗资源存储和传输结果。为确保资源释放,对每次调用任务后返回的每个 AsyncResult 实例,最终都必须调用 get() 或 forget()。
另请参阅
RPC结果后端(RabbitMQ/QPid)
RPC结果后端 rpc:// 比较特殊:它不实际存储状态,而是将状态作为消息发送。这意味着结果只能检索一次,而且只能由发起任务的客户端检索。两个不同进程无法等待同一结果。
尽管有这个限制,需要实时接收状态变化时,它仍是很好的选择。采用消息通信,客户端无需轮询新状态。
消息默认是临时的、不持久化的,所以消息代理重启后结果会消失。可以通过 result_persistent 设置,让结果后端发送持久化消息。
数据库结果后端
将状态保存在数据库中很方便,尤其是已经使用数据库的Web应用,但它也有局限。
-
轮询数据库中的新状态代价较高,应增大 result.get() 等操作的轮询间隔。
-
某些数据库默认的事务隔离级别不适合轮询表中的变化。
MySQL 默认隔离级别是 REPEATABLE-READ:当前事务提交之前,看不到其他事务所做的变化。
建议改为 READ-COMMITTED 隔离级别。
内置状态
PENDING
任务正在等待执行,或者任务未知。任何未知的任务ID都被视为处于等待状态。
SUCCESS
任务已成功执行。
- 元数据:
-
result 包含任务返回值。
- 传播异常:
-
是。
- 已就绪:
-
是。
FAILURE
任务执行失败。
- 元数据:
-
result 包含发生的异常,traceback 包含抛出异常时的堆栈回溯。
- 传播异常:
-
是。
RETRY
任务正在重试。
- 元数据:
-
result 包含引起重试的异常,traceback 包含抛出异常时的堆栈回溯。
- 传播异常:
-
否。
REVOKED
任务已被撤销。
- 传播异常:
-
是。
自定义状态
定义自定义状态很容易,只需要唯一名称,通常采用全大写字符串。例如,可查看定义了自定义 ABORTED 状态的abortable tasks。
使用 update_state() 更新任务状态:
@app.task(bind=True)
def upload_files(self, filenames):
for i, file in enumerate(filenames):
if not self.request.called_directly:
self.update_state(state='PROGRESS',
meta={'current': i, 'total': len(filenames)})
原作者在示例中创建了“PROGRESS”状态,告诉理解该状态的应用任务正在运行,并通过状态元数据中的 current 和 total 计数说明进展位置。这些信息可以用于创建进度条等界面。
创建可被pickle序列化的异常
Python 中一个不太为人所知的事实是:异常必须符合一些简单规则,才能被 pickle 模块序列化。
使用 Pickle 序列化器时,如果任务抛出的异常无法被 pickle 序列化,任务就不能正常工作。
要让异常支持 pickle,异常的 .args 属性必须包含创建实例时使用的原始参数。最简单的方法是在异常中调用 Exception.__init__。
下面给出一些可行示例,以及一个不可行示例:
# OK:
class HttpError(Exception):
pass
# BAD:
class HttpError(Exception):
def __init__(self, status_code):
self.status_code = status_code
# OK:
class HttpError(Exception):
def __init__(self, status_code):
self.status_code = status_code
Exception.__init__(self, status_code) # <-- REQUIRED
因此,规则是:对于支持自定义参数 *args 的异常,必须使用 Exception.__init__(self, *args)。
关键字参数没有特殊支持。如果希望在异常反序列化后保留关键字参数,必须将它们作为普通位置参数传递:
class HttpError(Exception):
def __init__(self, status_code, headers=None, body=None):
self.status_code = status_code
self.headers = headers
self.body = body
super(HttpError, self).__init__(status_code, headers, body)
半谓词
worker 将任务包裹在一个跟踪函数中,用来记录任务的最终状态。有些异常可以通知该函数,改变它对任务返回的处理方式。
Ignore
任务可以抛出 Ignore,强制 worker 忽略该任务。这意味着不记录任何任务状态,但消息仍会被确认并从队列删除。
可以用它实现自定义的类似撤销功能,或者手动保存任务结果。
将已撤销任务保存在 Redis 集合中的示例:
from celery.exceptions import Ignore
@app.task(bind=True)
def some_task(self):
if redis.ismember('tasks.revoked', self.request.id):
raise Ignore()
手动保存结果的示例:
from celery import states
from celery.exceptions import Ignore
@app.task(bind=True)
def get_tweets(self, user):
timeline = twitter.get_timeline(user)
if not self.request.called_directly:
self.update_state(state=states.SUCCESS, meta=timeline)
raise Ignore()
Reject
任务可以抛出 Reject,通过AMQP的 basic_reject 方法拒绝任务消息。只有启用 Task.acks_late 时,这才会生效。
拒绝消息与确认消息具有相同效果,但某些消息代理提供可利用的附加功能。例如,RabbitMQ 支持死信交换机,可以为队列配置死信交换机,将被拒绝的消息重新投递到那里。
Reject 也可以将消息重新入队,但必须谨慎使用,因为很容易形成无限消息循环。
任务导致内存不足时使用 reject 的示例:
import errno
from celery.exceptions import Reject
@app.task(bind=True, acks_late=True)
def render_scene(self, path):
file = get_file(path)
try:
renderer.render_scene(file)
# if the file is too big to fit in memory
# we reject it so that it's redelivered to the dead letter exchange
# and we can manually inspect the situation.
except MemoryError as exc:
raise Reject(exc, requeue=False)
except OSError as exc:
if exc.errno == errno.ENOMEM:
raise Reject(exc, requeue=False)
# For any other error we retry after 10 seconds.
except Exception as exc:
raise self.retry(exc, countdown=10)
将消息重新入队的示例:
from celery.exceptions import Reject
@app.task(bind=True, acks_late=True)
def requeues(self):
if not self.request.delivery_info['redelivered']:
raise Reject('no reason', requeue=True)
print('received two times')
basic_reject 方法的更多细节,参见所用消息代理的文档。
Retry
Task.retry 方法抛出 Retry,告诉 worker 任务正在重试。
自定义任务类
所有任务都继承自 app.Task,run() 方法成为任务函数体。
例如,以下代码:
@app.task
def add(x, y):
return x + y
在内部大致会转换成:
class _AddTask(app.Task):
def run(self, x, y):
return x + y
add = app.tasks[_AddTask.name]
实例化
不会为每个请求实例化任务,而是在任务注册表中将它注册为一个全局实例。
这意味着每个进程只调用一次 __init__ 构造函数,任务类在语义上更接近 Actor。
如果有如下任务:
from celery import Task
class NaiveAuthenticateServer(Task):
def __init__(self):
self.users = {'george': 'password'}
def run(self, username, password):
try:
return self.users[username] == password
except KeyError:
return False
并将所有请求路由到同一个进程,它就会在请求之间保留状态。
这也可用于缓存资源。例如,下面的任务基类缓存数据库连接:
from celery import Task
class DatabaseTask(Task):
_db = None
@property
def db(self):
if self._db is None:
self._db = Database.connect()
return self._db
用于单个任务
可以像这样应用到每个任务:
from celery.app import task
@app.task(base=DatabaseTask, bind=True)
def process_rows(self: task):
for row in self.db.table.all():
process_row(row)
之后,process_rows 任务的 db 属性在各进程内会始终保持相同。
用于整个应用
实例化 Celery 应用时,将自定义类作为 task_cls 参数传入,即可在整个应用中使用。参数可以是任务类本身,也可以是指定该类Python路径的字符串:
from celery import Celery
app = Celery('tasks', task_cls='your.module.path:DatabaseTask')
这会让应用中所有通过装饰器语法声明的任务使用 DatabaseTask 类,并都拥有 db 属性。
默认值是 Celery 提供的类 'celery.app.task:Task'。
处理方法
任务处理方法在生命周期的特定时刻执行。所有处理方法都在执行任务的同一 worker 进程和线程中同步运行。
执行时间线
下图展示确切执行顺序:
Worker Process Timeline
┌───────────────────────────────────────────────────────────────┐
│ 1. before_start() ← Blocks until complete │
│ 2. run() ← Your task function │
│ 3. [Result Backend] ← State + return value persisted │
│ 4. on_success() OR ← Outcome-specific handler │
│ on_retry() OR │ │
│ on_failure() │ │
│ 5. after_return() ← Runs last on terminal states │
│ (skipped for RETRY/REJECTED/IGNORED) │
└───────────────────────────────────────────────────────────────┘
重要
要点:
-
所有处理方法与任务本身运行在同一个 worker 进程中。
-
before_start会阻塞任务,在它结束之前run()不会开始。 -
结果后端在
on_success/on_failure之前更新,因此处理方法仍在运行时,其他客户端可能已经看到任务完成。 -
任务到达终态时才执行
after_return;对于RETRY、REJECTED或IGNORED不会执行。如果需要每次尝试都触发的hook,应使用task_postrun信号。
可用处理方法
- before_start(self, task_id, args, kwargs)
-
任务开始执行前由 worker 运行。
注意
此处理方法会阻塞任务:在
before_start返回之前,run()方法不会开始执行。5.2版本新增。
- 参数:
-
-
task_id:待执行任务的唯一ID。
-
args:待执行任务的原始位置参数。
-
kwargs:待执行任务的原始关键字参数。
-
此处理方法的返回值会被忽略。
- on_success(self, retval, task_id, args, kwargs)
-
成功处理方法。
任务成功执行时由 worker 运行。
注意
调用此方法前,任务结果已经持久化到结果后端。因此,此方法仍在执行时,外部客户端可能已经看到任务状态为
SUCCESS。- 参数:
-
-
retval:任务返回值。
-
task_id:已执行任务的唯一ID。
-
args:已执行任务的原始位置参数。
-
kwargs:已执行任务的原始关键字参数。
-
此处理方法的返回值会被忽略。
- on_retry(self, exc, task_id, args, kwargs, einfo)
-
重试处理方法。
任务准备重试时由 worker 运行。
注意
在结果后端已将任务状态更新为
RETRY、但重试尚未调度时调用。- 参数:
-
-
exc:传给
retry()的异常。 -
task_id:重试任务的唯一ID。
-
args:重试任务的原始位置参数。
-
kwargs:重试任务的原始关键字参数。
-
einfo:
ExceptionInfo实例。
-
此处理方法的返回值会被忽略。
- on_failure(self, exc, task_id, args, kwargs, einfo)
-
失败处理方法。
任务失败时由 worker 运行。
注意
调用此方法之前,任务结果已以
FAILURE状态持久化到结果后端。因此,此方法仍在执行时,外部客户端可能已经看到任务失败。- 参数:
-
-
exc:任务抛出的异常。
-
task_id:失败任务的唯一ID。
-
args:失败任务的原始位置参数。
-
kwargs:失败任务的原始关键字参数。
-
einfo:
ExceptionInfo实例。
-
此处理方法的返回值会被忽略。
- after_return(self, status, retval, task_id, args, kwargs, einfo)
-
任务返回后调用的处理方法。
注意
任务进入终态后,在对应结果的处理方法之后执行。
实际意味着在
on_success或on_failure之后执行;对于RETRY、REJECTED或IGNORED状态不会执行。如果需要每次尝试都执行的hook,可以使用task_postrun信号。- 参数:
-
-
status:当前任务状态。
-
retval:任务返回值或异常。
-
task_id:任务唯一ID。
-
args:已返回任务的原始位置参数。
-
kwargs:已返回任务的原始关键字参数。
-
einfo:
ExceptionInfo实例。
-
此处理方法的返回值会被忽略。
用法示例
import time
from celery import Task
class MyTask(Task):
def before_start(self, task_id, args, kwargs):
print(f"Task {task_id} starting with args {args}")
# This blocks - run() won't start until this returns
def on_success(self, retval, task_id, args, kwargs):
print(f"Task {task_id} succeeded with result: {retval}")
# Result is already visible to clients at this point
def on_failure(self, exc, task_id, args, kwargs, einfo):
print(f"Task {task_id} failed: {exc}")
# Task state is already FAILURE in backend
def after_return(self, status, retval, task_id, args, kwargs, einfo):
print(f"Task {task_id} finished with status: {status}")
# Always runs last
@app.task(base=MyTask)
def my_task(x, y):
return x + y
请求与自定义请求
收到执行任务的消息后,worker 会创建一个 request 表示这次执行请求。
自定义任务类可以修改 celery.app.task.Task.Request 属性,覆盖使用的请求类。可以直接赋值为自定义请求类,也可以赋值为它的完全限定名称。
请求承担多项职责,自定义请求类必须全部覆盖:它们负责实际运行并跟踪任务。强烈建议继承 celery.worker.request.Request。
使用 预派生进程worker 时,on_timeout() 和 on_failure() 在 worker 主进程中执行。应用可以据此发现 celery.app.task.Task.on_failure() 无法检测到的失败。
例如,以下自定义请求会检测并记录硬时间限制及其他失败:
import logging
from celery import Task
from celery.worker.request import Request
logger = logging.getLogger('my.package')
class MyRequest(Request):
'A minimal custom request to log failures and hard time limits.'
def on_timeout(self, soft, timeout):
super(MyRequest, self).on_timeout(soft, timeout)
if not soft:
logger.warning(
'A hard timeout was enforced for task %s',
self.task.name
)
def on_failure(self, exc_info, send_failed_event=True, return_ok=False):
super().on_failure(
exc_info,
send_failed_event=send_failed_event,
return_ok=return_ok
)
logger.warning(
'Failure detected for task %s',
self.task.name
)
class MyTask(Task):
Request = MyRequest # you can use a FQN 'my.package:MyRequest'
@app.task(base=MyTask)
def some_longrunning_task():
# use your imagination
工作原理
下面介绍技术细节。这部分并非必需知识,但可能有助于理解。
所有已定义任务都记录在注册表中,注册表包含任务名称及其任务类。可以自行查看:
>>> from proj.celery import app
>>> app.tasks
{'celery.chord_unlock':
<@task: celery.chord_unlock>,
'celery.backend_cleanup':
<@task: celery.backend_cleanup>,
'celery.chord':
<@task: celery.chord>}
上面列出的是 Celery 内置任务。注意,只有导入任务定义所在模块时,任务才会注册。
默认加载器会导入 imports 设置中列出的模块。
app.task() 装饰器负责将任务注册到应用的任务注册表。
发送任务时,不会发送实际函数代码,只发送待执行的任务名称。worker 收到消息后,在自己的任务注册表中查找名称,得到执行代码。
这意味着 worker 应始终更新为与客户端相同的软件版本。这是一个局限,而其他替代方案仍是尚未解决的技术挑战。
提示与最佳实践
忽略不需要的结果
如果不关心任务结果,务必设置 ignore_result,因为存储结果会浪费时间与资源。
@app.task(ignore_result=True)
def mytask():
something()
也可以通过 task_ignore_result 在全局禁用结果。
调用 apply_async 时传入布尔参数 ignore_result,可以按每次执行启用或禁用结果。
@app.task
def mytask(x, y):
return x + y
# No result will be stored
result = mytask.apply_async((1, 2), ignore_result=True)
print(result.get()) # -> None
# Result will be stored
result = mytask.apply_async((1, 2), ignore_result=False)
print(result.get()) # -> 3
配置结果后端后,任务默认不忽略结果,即 ignore_result=False。
选项优先级按以下顺序覆盖:
-
任务执行时的
ignore_result选项。
更多优化建议
更多优化建议见优化指南。
避免启动同步子任务
让一个任务等待另一个任务的结果效率很低,worker池耗尽时甚至会死锁。
应采用异步设计,例如使用回调。
不佳示例:
@app.task
def update_page_info(url):
page = fetch_page.delay(url).get()
info = parse_page.delay(page).get()
store_page_info.delay(url, info)
@app.task
def fetch_page(url):
return myhttplib.get(url)
@app.task
def parse_page(page):
return myparser.parse_document(page)
@app.task
def store_page_info(url, info):
return PageInfo.objects.create(url, info)
较好示例:
def update_page_info(url):
# fetch_page -> parse_page -> store_page
chain = fetch_page.s(url) | parse_page.s() | store_page_info.s(url)
chain()
@app.task()
def fetch_page(url):
return myhttplib.get(url)
@app.task()
def parse_page(page):
return myparser.parse_document(page)
@app.task(ignore_result=True)
def store_page_info(info, url):
PageInfo.objects.create(url=url, info=info)
原作者在这里将不同的 signature() 连接成任务链。关于链及其他强大组合方式,参见Canvas:设计工作流。
Celery 默认不允许在任务内部同步执行子任务,但在罕见或极端情况下可能确有需要。警告:不建议启用同步子任务执行。
@app.task
def update_page_info(url):
page = fetch_page.delay(url).get(disable_sync_subtasks=False)
info = parse_page.delay(page).get(disable_sync_subtasks=False)
store_page_info.delay(url, info)
@app.task
def fetch_page(url):
return myhttplib.get(url)
@app.task
def parse_page(page):
return myparser.parse_document(page)
@app.task
def store_page_info(url, info):
return PageInfo.objects.create(url, info)
性能与策略
粒度
任务粒度指每个子任务需要完成的计算量。通常,将问题拆成许多小任务,比少数长任务更好。
小任务可以提高并行处理数量,也不会运行过久而阻止 worker 处理其他等待任务。
但执行任务也有开销:需要发送消息,数据可能不在本地等。任务过细时,新增开销可能抵消全部收益。
另请参阅
Art of Concurrency一书有专门讨论任务粒度的章节[AOC1]。
Clay Breshears,《The Art of Concurrency》,第2.2.1节。O’Reilly Media, Inc.,2009年5月15日。ISBN-13:978-0-596-52153-0。
数据局部性
处理任务的 worker 应尽量靠近数据。最好是在内存中有一份副本,最差则是需要从另一个大洲完整传输数据。
数据距离很远时,可以尝试在数据所在地运行另一个 worker;如果做不到,可以缓存常用数据,或预加载已知即将使用的数据。
worker 之间共享数据最简单的方式是使用 memcached 这样的分布式缓存系统。
另请参阅
Jim Gray 的论文Distributed Computing Economics是数据局部性主题很好的入门读物。
状态
Celery 是分布式系统,无法预先知道任务会在哪个进程或哪台机器上执行,甚至无法确定任务能否及时运行。
异步编程有句话:“确认世界的状态,是任务的责任。”从发出任务请求到实际执行,外部状态可能已经改变,因此任务必须自行确认状态符合要求。例如,任务为搜索引擎重建索引,而索引最多每5分钟重建一次,那么检查这一条件应由任务负责,而不是由调用方负责。
另一个陷阱是 Django 模型对象:不应直接将它们作为任务参数。通常更好的办法是在任务执行时从数据库重新获取对象,因为使用旧数据可能导致竞态条件。
设想一个场景:有一篇文章,以及一个自动展开其中缩写的任务。
class Article(models.Model):
title = models.CharField()
body = models.TextField()
@app.task
def expand_abbreviations(article):
article.body.replace('MyCorp', 'My Corporation')
article.save()
作者先创建并保存文章,再点击按钮启动缩写处理任务:
>>> article = Article.objects.get(id=102)
>>> expand_abbreviations.delay(article)
队列繁忙,任务要2分钟后才会执行。在此期间,另一位作者修改了文章。任务终于运行时,由于参数里保存的是旧正文,会把文章恢复成旧版本。
解决这个竞态条件很简单:只传文章ID,并在任务函数体中重新获取文章:
@app.task
def expand_abbreviations(article_id):
article = Article.objects.get(id=article_id)
article.body.replace('MyCorp', 'My Corporation')
article.save()
>>> expand_abbreviations.delay(article_id)
这样也可能提高性能,因为发送大消息的成本可能很高。
数据库事务
再看另一个例子:
from django.db import transaction
from django.http import HttpResponseRedirect
@transaction.atomic
def create_article(request):
article = Article.objects.create()
expand_abbreviations.delay(article.pk)
return HttpResponseRedirect('/articles/')
这个 Django 视图在数据库中创建文章对象,再把主键传给任务。它使用 transaction.atomic 装饰器:视图返回时提交事务,视图抛出异常时回滚。
事务的原子性在这里导致竞态条件:视图函数返回响应后,文章对象才会持久化到数据库。如果异步任务在事务提交之前开始执行,就可能在对象尚不存在时查询它。因此,需要确保事务提交之后才触发任务。
解决方法是改用 delay_on_commit():
from django.db import transaction
from django.http import HttpResponseRedirect
@transaction.atomic
def create_article(request):
article = Article.objects.create()
expand_abbreviations.delay_on_commit(article.pk)
return HttpResponseRedirect('/articles/')
该方法在 Celery 5.4 中加入,是对 Django on_commit 回调的快捷封装,会在所有事务成功提交后启动 Celery 任务。
Celery 5.4之前的版本
使用较旧版本时,可以直接使用 Django 回调实现相同效果:
import functools
from django.db import transaction
from django.http import HttpResponseRedirect
@transaction.atomic
def create_article(request):
article = Article.objects.create()
transaction.on_commit(
functools.partial(expand_abbreviations.delay, article.pk)
)
return HttpResponseRedirect('/articles/')
注意
Django 1.9及更新版本提供 on_commit。使用更旧版本时,可以通过 django-transaction-hooks 库增加支持。
示例
来看实际场景:博客需要过滤评论中的垃圾信息。评论创建后,在后台运行垃圾信息过滤,用户无需等待过滤完成。
原作者的 Django 博客应用允许用户评论文章。以下介绍该应用部分模型、视图和任务。
blog/models.py
评论模型如下:
from django.db import models
from django.utils.translation import ugettext_lazy as _
class Comment(models.Model):
name = models.CharField(_('name'), max_length=64)
email_address = models.EmailField(_('email address'))
homepage = models.URLField(_('home page'),
blank=True, verify_exists=False)
comment = models.TextField(_('comment'))
pub_date = models.DateTimeField(_('Published date'),
editable=False, auto_add_now=True)
is_spam = models.BooleanField(_('spam?'),
default=False, editable=False)
class Meta:
verbose_name = _('comment')
verbose_name_plural = _('comments')
在处理评论提交的视图中,原作者先将评论写入数据库,然后在后台启动垃圾信息过滤任务。
blog/views.py
from django import forms
from django.http import HttpResponseRedirect
from django.template.context import RequestContext
from django.shortcuts import get_object_or_404, render_to_response
from blog import tasks
from blog.models import Comment
class CommentForm(forms.ModelForm):
class Meta:
model = Comment
def add_comment(request, slug, template_name='comments/create.html'):
post = get_object_or_404(Entry, slug=slug)
remote_addr = request.META.get('REMOTE_ADDR')
if request.method == 'post':
form = CommentForm(request.POST, request.FILES)
if form.is_valid():
comment = form.save()
# Check spam asynchronously.
tasks.spam_filter.delay(comment_id=comment.id,
remote_addr=remote_addr)
return HttpResponseRedirect(post.get_absolute_url())
else:
form = CommentForm()
context = RequestContext(request, {'form': form})
return render_to_response(template_name, context_instance=context)
原作者使用 Akismet 过滤垃圾评论;WordPress 免费博客平台也使用这项服务。原文说明,Akismet 可供个人免费使用,商业使用需要付费,注册服务后才能取得API密钥。
原作者使用 Michael Foord 编写的 akismet.py 库调用 Akismet API。
blog/tasks.py
from celery import Celery
from akismet import Akismet
from django.core.exceptions import ImproperlyConfigured
from django.contrib.sites.models import Site
from blog.models import Comment
app = Celery(broker='amqp://')
@app.task
def spam_filter(comment_id, remote_addr=None):
logger = spam_filter.get_logger()
logger.info('Running spam filter for comment %s', comment_id)
comment = Comment.objects.get(pk=comment_id)
current_domain = Site.objects.get_current().domain
akismet = Akismet(settings.AKISMET_KEY, 'http://{0}'.format(domain))
if not akismet.verify_key():
raise ImproperlyConfigured('Invalid AKISMET_KEY')
is_spam = akismet.comment_check(user_ip=remote_addr,
comment_content=comment.comment,
comment_author=comment.name,
comment_author_email=comment.email_address)
if is_spam:
comment.is_spam = True
comment.save()
return is_spam











暂无评论内容