本文依据 Celery 5.6 稳定版文档,介绍自动与手动路由、消息优先级、AMQP 基础以及自定义路由器和广播。开发版文档见开发分支。
注意:并非所有传输后端都支持 topic、fanout 等路由方式,选择前应参考 Celery 的传输后端比较表。
基础用法
自动路由
最简单的路由方式是使用 task_create_missing_queues,该选项默认开启。开启后,如果指定的队列尚未在 task_queues 中定义,Celery 会自动创建它,适合简单的任务分流。
假设服务器 x、y 处理普通任务,而服务器 z 只处理 feed 相关任务,可以这样配置:
task_routes = {'feed.tasks.import_feed': {'queue': 'feeds'}}
启用后,导入 feed 的任务会发送到 feeds 队列,其他任务进入默认队列。出于历史原因,默认队列名为 celery。
也可以使用 glob 通配符或正则表达式,匹配 feed.tasks 命名空间中的所有任务:
app.conf.task_routes = {'feed.tasks.*': {'queue': 'feeds'}}
如果匹配顺序很重要,应改用条目列表形式定义路由:
task_routes = ([
('feed.tasks.*', {'queue': 'feeds'}),
('web.tasks.*', {'queue': 'web'}),
(re.compile(r'(video|image)\.tasks\..*'), {'queue': 'media'}),
],)
注意:task_routes 可以是字典,也可以是路由器对象列表。因此这里需要用一个包含列表的元组来表达有序匹配。
安装路由配置后,可以让服务器 z 只处理 feeds 队列:
user@z:/$ celery -A proj worker -Q feeds
可以指定任意数量的队列,因此也可以让它同时处理默认队列:
user@z:/$ celery -A proj worker -Q feeds,celery
修改默认队列名称
通过以下配置修改默认队列名:
app.conf.task_default_queue = 'default'
自动队列是如何定义的
自动路由的目的,是为只有基础需求的用户隐藏复杂的 AMQP 配置。例如,名为 video 的队列会使用以下设置创建:
{'exchange': 'video',
'exchange_type': 'direct',
'routing_key': 'video'}
Redis、SQS 等非 AMQP 后端不支持交换机,因此要求交换机与队列名称相同。这种设计使自动队列也能在这些后端上工作。
手动路由
仍以 x、y 处理普通任务、z 只处理 feed 为例,可以显式定义队列:
from kombu import Queue
app.conf.task_default_queue = 'default'
app.conf.task_queues = (
Queue('default', routing_key='task.#'),
Queue('feed_tasks', routing_key='feed.#'),
)
app.conf.task_default_exchange = 'tasks'
app.conf.task_default_exchange_type = 'topic'
app.conf.task_default_routing_key = 'task.default'
task_queues 是 Queue 实例组成的列表。某个队列未指定交换机或交换机类型时,会采用 task_default_exchange 与 task_default_exchange_type。
要把任务发往 feed_tasks,可在 task_routes 中添加规则:
task_routes = {
'feeds.tasks.import_feed': {
'queue': 'feed_tasks',
'routing_key': 'feed.import',
},
}
也可以在 Task.apply_async() 或 send_task() 中传入 routing_key 等路由参数,覆盖配置:
>>> from feeds.tasks import import_feed
>>> import_feed.apply_async(args=['http://cnn.com/rss'],
... queue='feed_tasks',
... routing_key='feed.import')
通过 celery worker -Q 让服务器 z 只消费 feed 队列:
user@z:/$ celery -A proj worker -Q feed_tasks --hostname=z@%h
服务器 x、y 应配置为消费默认队列:
user@x:/$ celery -A proj worker -Q default --hostname=x@%h
user@y:/$ celery -A proj worker -Q default --hostname=y@%h
在工作量较大时,也可以让处理 feed 的 worker 同时接收普通任务:
user@z:/$ celery -A proj worker -Q feed_tasks,default --hostname=z@%h
如果要添加位于另一个交换机上的队列,明确指定自定义交换机和交换机类型即可:
from kombu import Exchange, Queue
app.conf.task_queues = (
Queue('feed_tasks', routing_key='feed.#'),
Queue('regular_tasks', routing_key='task.#'),
Queue('image_tasks', exchange=Exchange('mediatasks', type='direct'),
routing_key='image.compress'),
)
如果不熟悉这些术语,可以先阅读下文的 AMQP 入门。原文还推荐有关队列与交换机的文章、CloudAMQP 教程和 RabbitMQ FAQ 作为补充材料。
特殊路由选项
RabbitMQ 消息优先级
支持的传输后端:RabbitMQ。此功能自 Celery 4.0 加入。
通过队列参数 x-max-priority 开启优先级支持:
from kombu import Exchange, Queue
app.conf.task_queues = [
Queue('tasks', Exchange('tasks'), routing_key='tasks',
queue_arguments={'x-max-priority': 10}),
]
可以用 task_queue_max_priority 为所有队列设置默认值:
app.conf.task_queue_max_priority = 10
也可以用 task_default_priority 为所有任务设置默认优先级:
app.conf.task_default_priority = 5
Redis 消息优先级
支持的传输后端:Redis。
Celery 的 Redis 传输层会处理优先级字段,但 Redis 本身没有消息优先级概念。因此在使用前应了解其实现方式,避免出现与预期不同的行为。
要按优先级调度任务,配置 queue_order_strategy 传输选项:
app.conf.broker_transport_options = {
'queue_order_strategy': 'priority',
}
其实现方式是为每个逻辑队列创建多个 Redis 列表。虽然优先级有 0~9 共十个级别,但默认会合并为四档以节约资源。因此,逻辑上名为 celery 的队列,实际会拆分为四个列表。
最高优先级的队列仍名为 celery,其余名称会附加分隔符和优先级数字。默认分隔符为 \x06\x16:
['celery', 'celery\x06\x163', 'celery\x06\x166', 'celery\x06\x169']
如需更多优先级档位或不同分隔符,设置 priority_steps 与 sep:
app.conf.broker_transport_options = {
'priority_steps': list(range(10)),
'sep': ':',
'queue_order_strategy': 'priority',
}
以上配置会生成以下队列名:
['celery', 'celery:1', 'celery:2', 'celery:3', 'celery:4', 'celery:5', 'celery:6', 'celery:7', 'celery:8', 'celery:9']
这种机制无法完全等同于消息代理服务器原生实现的优先级,最多是近似效果,但对某些应用仍然足够。
AMQP 入门
消息
消息由头部与正文组成。Celery 在头部保存内容类型和内容编码。内容类型通常对应消息的序列化格式。正文包含要执行的任务名称、任务 ID(UUID)、调用参数,以及重试次数、预计执行时间 ETA 等元数据。
用 Python 字典表示的任务消息示例:
{'task': 'myapp.tasks.add',
'id': '54086c5e-6193-4575-8308-dbab76798756',
'args': [4, 4],
'kwargs': {}}
生产者、消费者与消息代理
发送消息的客户端通常称为发布者或生产者,接收消息的一方称为消费者。消息代理(broker)是消息服务器,负责把生产者发送的消息路由给消费者。这些术语会频繁出现在 AMQP 资料中。
交换机、队列与路由键
消息的基本流转过程如下:
- 消息首先发送到交换机。
- 交换机将消息路由到一个或多个队列。不同交换机类型提供不同路由方式,适合不同消息场景。
- 消息在队列中等待消费。
- 消费者确认消息后,它才从队列中删除。
发送和接收消息前,需要创建交换机、创建队列,并把队列绑定到交换机。
Celery 会自动创建 task_queues 所需的实体,除非队列的 auto_declare 被设为 False。下面定义三个队列,分别处理视频、图片和其他默认任务:
from kombu import Exchange, Queue
app.conf.task_queues = (
Queue('default', Exchange('default'), routing_key='default'),
Queue('videos', Exchange('media'), routing_key='media.video'),
Queue('images', Exchange('media'), routing_key='media.image'),
)
app.conf.task_default_queue = 'default'
app.conf.task_default_exchange_type = 'direct'
app.conf.task_default_routing_key = 'default'
交换机类型
交换机类型决定消息如何路由。AMQP 标准定义了 direct、topic、fanout 和 headers。RabbitMQ 还可通过插件提供非标准类型,例如 Michael Bridgen 的 last-value-cache 插件。
Direct 交换机
Direct 交换机精确匹配路由键。使用 video 路由键绑定的队列,只接收路由键恰好为 video 的消息。
Topic 交换机
Topic 交换机按点分隔的词匹配路由键,支持两个通配符:* 匹配一个词,# 匹配零个或多个词。
例如,有 usa.news、usa.weather、norway.news、norway.weather 四个路由键时:
*.news匹配所有新闻。usa.#匹配所有美国相关消息。usa.weather只匹配美国天气。
相关 API 命令
exchange.declare(exchange_name, type, passive, durable, auto_delete, internal) 按名称声明交换机,对应 Channel.exchange_declare。
passive:不创建交换机,只检查它是否已存在。durable:交换机持久化,可在消息代理重启后保留。auto_delete:当没有队列使用该交换机时,由消息代理自动删除。
queue.declare(queue_name, passive, durable, exclusive, auto_delete) 按名称声明队列,对应 Channel.queue_declare。独占队列只能由当前连接消费;exclusive 还隐含 auto_delete。
queue.bind(queue_name, exchange_name, routing_key) 通过路由键将队列绑定到交换机,对应 Channel.queue_bind。未绑定的队列收不到消息,因此这一步是必需的。
queue.delete(name, if_unused=False, if_empty=False) 删除队列及其绑定,对应 Channel.queue_delete。
exchange.delete(name, if_unused=False) 删除交换机,对应 Channel.exchange_delete。
注意:“声明”不一定表示“创建”,而是确认实体存在且可用。协议不规定由生产者还是消费者首次创建交换机、队列或绑定,通常由最先需要它的一方创建。
动手使用 API
Celery 附带 celery amqp 工具,可以从命令行访问 AMQP API,执行创建和删除队列或交换机、清空队列、发送消息等管理操作。它也可用于非 AMQP 消息代理,但具体实现未必支持全部命令。
既可以直接把命令作为 celery amqp 参数传入,也可以不带参数启动交互式 shell:
$ celery -A proj amqp
-> connecting to amqp://guest@localhost:5672/.
-> connected.
1>
1> 是提示符,数字表示已执行命令的计数。输入 help 可查看可用命令。工具也支持自动补全,输入命令前缀后按 Tab 可显示匹配项。
先创建一个可发送消息的队列:
$ celery -A proj amqp
1> exchange.declare testexchange direct
ok.
2> queue.declare testqueue
ok. queue:testqueue messages:0 consumers:0.
3> queue.bind testqueue testexchange testkey
ok.
这里创建了 direct 类型的 testexchange 交换机和名为 testqueue 的队列,并用 testkey 路由键建立绑定。
之后,发送到 testexchange 且路由键为 testkey 的消息都会进入该队列。使用 basic.publish 发送消息:
4> basic.publish 'This is a message!' testexchange testkey
ok.
消息发送后,可以用 basic.get 同步轮询队列并取出消息。维护任务可以这样使用,持续运行的服务则应采用 basic.consume。
5> basic.get testqueue
{'body': 'This is a message!',
'delivery_info': {'delivery_tag': 1,
'exchange': u'testexchange',
'message_count': 0,
'redelivered': False,
'routing_key': u'testkey'},
'properties': {}}
AMQP 使用确认机制表示消息已成功接收和处理。如果消息尚未确认,而消费者通道关闭,该消息会重新交给其他消费者。
注意输出中的 delivery tag。每条消息在同一连接通道内都有唯一的交付标签,用于确认该消息。但标签并不跨连接唯一,因此另一个客户端中的标签 1 可能指向完全不同的消息。
使用 basic.ack 确认收到的消息:
6> basic.ack 1
ok.
测试会话结束后,删除先前创建的实体:
7> queue.delete testqueue
ok. 0 messages deleted.
8> exchange.delete testexchange
ok.
任务路由详解
定义队列
Celery 通过 task_queues 定义可用队列。下面仍以视频、图片、默认任务三个队列为例:
default_exchange = Exchange('default', type='direct')
media_exchange = Exchange('media', type='direct')
app.conf.task_queues = (
Queue('default', default_exchange, routing_key='default'),
Queue('videos', media_exchange, routing_key='media.video'),
Queue('images', media_exchange, routing_key='media.image')
)
app.conf.task_default_queue = 'default'
app.conf.task_default_exchange = 'default'
app.conf.task_default_routing_key = 'default'
未显式匹配路由的任务会使用 task_default_queue。默认交换机、交换机类型和路由键,既是任务的默认路由参数,也为 task_queues 中未指定相应值的条目提供默认值。
同一队列还支持多个绑定。例如,两个不同路由键可以绑定到同一个队列:
from kombu import Exchange, Queue, binding
media_exchange = Exchange('media', type='direct')
CELERY_QUEUES = (
Queue('media', [
binding(media_exchange, routing_key='media.video'),
binding(media_exchange, routing_key='media.image'),
]),
)
指定任务目标
任务目标依次由以下来源决定:
- 传给
Task.apply_async()的路由参数。 - Task 自身定义的路由相关属性。
task_routes中的路由器。
通常建议避免把路由设置硬编码到业务代码中,而是通过路由器作为配置维护,这样最灵活。仍然可以在任务属性中设置合理默认值。
路由器
路由器是决定任务路由选项的函数。定义一个参数签名为 (name, args, kwargs, options, task=None, **kw) 的函数即可:
def route_task(name, args, kwargs, options, task=None, **kw):
if name == 'myapp.tasks.compress_video':
return {'exchange': 'video',
'exchange_type': 'topic',
'routing_key': 'video.compress'}
如果返回值包含 queue,Celery 会根据 task_queues 中的队列定义展开完整设置。例如:
{'queue': 'video', 'routing_key': 'video.compress'}
会展开为:
{'queue': 'video',
'exchange': 'video',
'exchange_type': 'topic',
'routing_key': 'video.compress'}
在 task_routes 中添加路由器类即可安装:
task_routes = (route_task,)
也可以按名称引用路由函数:
task_routes = ('myapp.routers.route_task',)
对于简单的“任务名到路由”映射,直接在 task_routes 中放入字典即可获得相同行为:
task_routes = {
'myapp.tasks.compress_video': {
'queue': 'video',
'routing_key': 'video.compress',
},
}
Celery 会按顺序遍历路由器,遇到第一个返回真值的路由器就停止,并将其结果作为最终路由。
可以用序列配置多个路由器:
task_routes = [
route_task,
{
'myapp.tasks.compress_video': {
'queue': 'video',
'routing_key': 'video.compress',
},
]
此时同样按顺序访问,选择第一个返回有效值的结果。
使用 Redis 或 RabbitMQ 时,还可以在路由中指定默认优先级:
task_routes = {
'myapp.tasks.compress_video': {
'queue': 'video',
'routing_key': 'video.compress',
'priority': 10,
},
}
任务调用 apply_async 时指定的优先级会覆盖该默认值:
task.apply_async(priority=0)
优先级顺序与集群响应能力
worker 会预取任务,因此同时提交大量任务时,最初的执行顺序可能不完全符合优先级。禁用预取可以避免这一情况,但会降低短小、快速任务的执行效率。
多数情况下,将 worker_prefetch_multiplier 降至 1,是一种更简单的折中方式:提高系统响应能力,同时避免完全禁用预取的代价。
注意,Redis 消息代理的优先级数值顺序相反,0 表示最高优先级。
广播
Celery 也支持广播路由。下面的 broadcast_tasks 交换机会把任务副本发送给所有连接到它的 worker:
from kombu.common import Broadcast
app.conf.task_queues = (Broadcast('broadcast_tasks'),)
app.conf.task_routes = {
'tasks.reload_cache': {
'queue': 'broadcast_tasks',
'exchange': 'broadcast_tasks'
}
}
这样,tasks.reload_cache 会发送给每个消费该广播队列的 worker。
也可以配合 celery beat 定时计划使用广播:
from kombu.common import Broadcast
from celery.schedules import crontab
app.conf.task_queues = (Broadcast('broadcast_tasks'),)
app.conf.beat_schedule = {
'test-task': {
'task': 'tasks.reload_cache',
'schedule': crontab(minute=0, hour='*/3'),
'options': {'exchange': 'broadcast_tasks'}
},
}
广播与结果
Celery 没有定义两个任务拥有相同 task_id 时结果后端应如何处理。如果同一任务被分发给多个 worker,其状态历史可能无法完整保存。因此,这类任务通常应设置 task.ignore_result。
原文来源:Routing Tasks。本文依据留存原文译为中文,代码示例按原文保留。
原文版权声明:Copyright。











暂无评论内容