Celery 任务路由指南

本文依据 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 资料中。

交换机、队列与路由键

消息的基本流转过程如下:

  1. 消息首先发送到交换机。
  2. 交换机将消息路由到一个或多个队列。不同交换机类型提供不同路由方式,适合不同消息场景。
  3. 消息在队列中等待消费。
  4. 消费者确认消息后,它才从队列中删除。

发送和接收消息前,需要创建交换机、创建队列,并把队列绑定到交换机。

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'),
    ]),
)

指定任务目标

任务目标依次由以下来源决定:

  1. 传给 Task.apply_async() 的路由参数。
  2. Task 自身定义的路由相关属性。
  3. 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。

© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容