基于 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
现在开始开发周期:
tox # 以确保测试通过。
在tests/中编写新测试并确保它们失败。
在channels_rabbitmq/中编写新代码以使测试通过。
提交拉取请求。
部署
使用semver。
git push并确保 Travis 测试全部通过。
git 标签 vX.XX
git push --标签
TravisCI 将推送到 PyPi。
项目详情
下载文件
下载适用于您平台的文件。如果您不确定要选择哪个,请了解有关安装包的更多信息。