Skip to main content

基于 RabbitMQ 的 ASGI 通道层实现

项目描述

使用 RabbitMQ 作为其后备存储的 Django Channels 通道层。

不支持Worker 和 Background Tasks。(请参阅基本原理 并使用await get_channel_layer().current_connection发送到作业队列。)

适用于 Python 3.8 或 3.9。

安装

点安装频道_rabbitmq

用法

然后在 Django 设置文件中设置通道层,如下所示:

CHANNEL_LAYERS = {
    "default": {
        "BACKEND": "channels_rabbitmq.core.RabbitmqChannelLayer",
        "CONFIG": {
            "host": "amqp://guest:guest@127.0.0.1/asgi",
            # "ssl_context": ... (optional)
        },
    },
}

下面列出了CONFIG的可能选项。

主持人

要连接的服务器的 URL,遵循RabbitMQ 规范。要连接到 RabbitMQ 集群,请使用 DNS 服务器将主机名解析为多个 IP 地址。如果在断开连接的情况下至少可以访问其中一个,channels_rabbitmq 将自动重新连接。

到期

一条消息在被静默丢弃之前应该在 RabbitMQ 队列中等待的最小秒数。

默认为60。您通常不需要更改此设置,但如果您希望减少高峰流量,您可能希望将其关闭,或者如果您希望积压高峰流量,则将其调高,直到您到达为止。

本地容量

在内存中排队的传入消息数。默认为100。(发送到具有两个通道的组的消息计为一条消息。)当local_capacity 消息排队时,RabbitMQ 上的消息积压将增加。

(这控制了RabbitMQ 队列上的prefetch_count 。)

local_expiry

从 RabbitMQ 接收到的消息必须在内存中等待receive()的最小秒数,然后才能删除它。默认为 到期

当消息在本地过期时,将记录一个警告。警告可以表明一个通道的消息比它可以处理的多;或者消息正在发送到不存在的通道。(也许group_add()暗示了缺少的频道,并且 从未调用过匹配的group_discard() 。)

如果local_expiry < expiry,那么当它们仍然存在于 RabbitMQ 队列中时,您最终可以在本地忽略(和记录)消息。这些消息将被确认,因此 RabbitMQ 的行为就像它们已被传递一样。

远程容量

每个客户端存储在 RabbitMQ 上的消息数。默认为100。(发送到两个不同客户端上具有三个通道的组的消息算作两条消息。)当remote_capacity消息在 RabbitMQ 中排队时,通道将拒绝新消息。从任何客户端调用send()group_send()到满容量客户端将引发ChannelFull

ssl_context

SSL上下文。将默认主机端口更改为 5671(而不是 5672)。

例如,要连接到将验证您的客户端的 TLS RabbitMQ 服务:

import ssl
ssl_context = ssl.create_default_context(
    cafile=str(Path(__file__).parent.parent / 'ssl' / 'server.cert'),
)
ssl_context.load_cert_chain(
    certfile=str(Path(__file__).parent.parent / 'ssl' / 'client.certchain'),
    keyfile=str(Path(__file__).parent.parent / 'ssl' / 'client.key'),
)
CHANNEL_LAYERS['default']['CONFIG']['ssl_context'] = ssl_context

默认情况下,没有 SSL 上下文;所有消息(和密码)都以明文形式传输。

groups_exchange

通道用于交换组消息的全局直接交换名称。默认为“组”。另请参阅设计决策

访问 Carehare

我们使用carehare来彻底处理错误。

Django Channels 的规范没有考虑“连接”和“断开”。这一层通过不断地重新连接,永远做到最好。

调用await get_channel_layer().current_connection以访问打开的 Carehare 连接。这使您可以使用没有“工作人员和后台任务”的作业队列。像这样:

# raise asyncio.CancelledError on failure
connection = await get_channel_layer().carehare_connection

# raise carehare.ConnectionClosed or carehare.ChannelClosed on error
await connection.publish(b"task", routing_key="job_queue")

(Carehare 文档解释了如何培养工人。)

关于错误的说明:“已连接”连接不能保证在每次发布期间都保持连接。它只是在过去的某个时间点连接的。当发生断开连接时,该连接上的所有挂起操作都会引发 carehare.ConnectionClosed。这个通道层将记录错误, get_channel_layer().carehare_connection将指向一个新的 Future。(这个错误+重新连接保证在生产中发生。)

设计决策

为了极大地扩展,这一层只为每个实例创建一个 RabbitMQ 队列。这意味着无论打开多少个 websocket 连接,一台 Web 服务器都会获得一个 RabbitMQ 队列。对于发送的每条消息,客户端层确定 RabbitMQ 队列名称并将其用作路由键。

默认情况下,组是使用称为“组”的单个全局 RabbitMQ 直接交换实现的。要将消息发送到组,该层将消息发送到“组”交换,组名称作为路由键。客户端在group_add()group_remove()期间绑定和取消绑定,以确保其任何组的消息都能到达它。另请参见groups_exchange 选项。

RabbitMQ 队列是独占的:当客户端断开连接(通过关闭或崩溃)时,RabbitMQ 将删除队列并取消绑定组。

创建连接后,它会污染事件循环,因此 如果在async_to_sync()中创建了连接,则async_to_sync()将破坏该连接 。每个连接都会启动一个后台异步循环,从 RabbitMQ 拉取消息并将它们路由到接收者队列;每个receive() 查询接收者队列。没有连接的空队列被删除。

与通道层规范的偏差

通道层规范屈服于 Redis 相关的限制。RabbitMQ 无法模拟 Redis。以下是不同之处:

  • 没有 ``flush`` 扩展:要刷新所有状态,只需断开所有客户端。(RabbitMQ 不允许一个客户端删除另一个客户端的数据结构。)

  • 没有 ``group_expiry`` 选项: 当group_add()没有匹配的group_discard()时, group_expiry 选项恢复。但是“组成员到期”的逻辑有一个致命的缺陷:它断开了合法成员的连接。channels_rabbitmq解决了每个根本问题:

    • Web 服务器崩溃:当 Web 服务器断开连接时,RabbitMQ 会擦除与 Web 服务器相关的所有状态。group_expiry解决这里没有问题 。

    • 编程错误:您可能会犯错并调用group_add()而不最终调用group_discard()。Redis 无法检测到此编程错误(因为它无法检测到 Web 服务器崩溃)。RabbitMQ 可以。local_expiry选项在您错误地错过group_discard()后使您的站点保持运行。丢弃过期消息时,通道层会发出警告。监控您的服务器日志以检测您的错误。

  • 没有“正常通道”正常通道 是作业队列。在大多数项目中,“正常通道”阅读器是工作进程,理想情况下与 Websockets 和 Django 分离。

    如果你想要一个异步的、基于 RabbitMQ 的作业队列,请研究carehare

依赖项

您需要 Python 3.8+ 和 RabbitMQ 服务器。

如果您有 Docker,以下是启动开发服务器的方法:

ssl/prepare-certs.sh  # Create SSL certificates used in tests
docker run --rm -it \
     -p 5671:5671 \
     -p 5672:5672 \
     -p 15672:15672 \
     -v "/$(pwd)"/ssl:/ssl \
     -e RABBITMQ_SSL_CACERTFILE=/ssl/ca.cert \
     -e RABBITMQ_SSL_CERTFILE=/ssl/server.cert \
     -e RABBITMQ_SSL_KEYFILE=/ssl/server.key \
     -e RABBITMQ_SSL_VERIFY=verify_peer \
     -e RABBITMQ_SSL_FAIL_IF_NO_PEER_CERT=true \
     rabbitmq:3.7.8-management-alpine

您可以通过http://localhost:15672访问 RabbitMQ 管理界面。

贡献

添加功能和修复错误

首先,启动一个开发 RabbitMQ 服务器:

ssl/prepare-certs.sh  # Create SSL certificates used in tests
docker run --rm -it \
     -p 5671:5671 \
     -p 5672:5672 \
     -p 15672:15672 \
     -v "/$(pwd)"/ssl:/ssl \
     -e RABBITMQ_SSL_CACERTFILE=/ssl/ca.cert \
     -e RABBITMQ_SSL_CERTFILE=/ssl/server.cert \
     -e RABBITMQ_SSL_KEYFILE=/ssl/server.key \
     -e RABBITMQ_SSL_VERIFY=verify_peer \
     -e RABBITMQ_SSL_FAIL_IF_NO_PEER_CERT=true \
     rabbitmq:3.8.11-management-alpine

现在开始开发周期:

  1. tox # 以确保测试通过。

  2. 在tests/中编写新测试并确保它们失败。

  3. 在channels_rabbitmq/中编写新代码以使测试通过。

  4. 提交拉取请求。

部署

使用semver

  1. git push并确保 Travis 测试全部通过。

  2. git 标签 vX.XX

  3. git push --标签

TravisCI 将推送到 PyPi。

项目详情


下载文件

下载适用于您平台的文件。如果您不确定要选择哪个,请了解有关安装包的更多信息。

源分布

channels_rabbitmq-4.0.0.tar.gz (19.8 kB 查看哈希

已上传 source

内置分布

channels_rabbitmq-4.0.0-py3-none-any.whl (17.0 kB 查看哈希

已上传 py3