- 后端
- WebSocket
- 异步编程
【免费下载链接】channels
Developer-friendly asynchrony for Django
导读
本指南围绕 Channels 中最核心的编程抽象——**Consumers(消费者)**展开。Channels 建立在低层 ASGI 规范之上,而 ASGI 本身更强调互操作性而非复杂的业务开发体验,因此 Channels 提供了 Consumers 这一高层抽象,让你能以事件回调函数的方式快速构建 ASGI 应用。读完本文,你将掌握基础消费者(AsyncConsumer/SyncConsumer)的事件分发机制、self.send与scope的用法、连接关闭处理,以及四类通用消费者(WebSocket/HTTP 系列的同步与异步版本)的完整实战用法,并能用仓库内的测试代码验证自己的实现。
一、为什么需要 Consumers:ASGI 之上的高层抽象
Channels 围绕一个低层规范——ASGI(文档见 docs/asgi.rst)构建。ASGI 设计的首要目标是让不同协议服务器、框架之间可以互操作,而不是让人用它直接写出复杂业务应用。Consumers 正是 Channels 为补齐这一短板而提供的丰富抽象。
Consumers 主要解决两件事:
- 把代码结构化为"事件发生即调用"的一组函数,你不需要自己手写事件循环(event loop);
- 允许你编写同步或异步代码,并替你处理两者之间的交接(handoff)与线程调度。
当然,你完全可以不用 Consumers,把 Channels 的其余部分——路由(docs/topics/routing.rst)、会话处理、认证——配合任意 ASGI 应用使用,但 Consumers 通常是编写应用代码的最佳方式。
二、基础布局:AsyncConsumer 与 SyncConsumer
一个 Consumer 是channels.consumer.AsyncConsumer或channels.consumer.SyncConsumer的子类。如名字所示:前者期望你编写异步代码,后者会把你的代码放到线程池(threadpool)里同步运行。
2.1 SyncConsumer 的最小示例:WebSocket 回声服务器
from channels.consumer import SyncConsumer class EchoConsumer(SyncConsumer): def websocket_connect(self, event): self.send({ "type": "websocket.accept", }) def websocket_receive(self, event): self.send({ "type": "websocket.send", "text": event["text"], })这是一个非常简单的 WebSocket 回声服务器:接受所有传入的 WebSocket 连接,然后把收到的每个文本帧原样回给客户端。
2.2 事件类型到方法名的映射规则
Consumers 围绕一系列命名方法组织,方法名由消息的type值推导而来:把type中的.全部替换为_。上面两个处理器分别对应websocket.connect与websocket.receive消息。这一机制在源码 channels/consumer.py 的get_handler_name中实现:
def get_handler_name(message): if "type" not in message: raise ValueError("Incoming message has no 'type' attribute") handler_name = message["type"].replace(".", "_") if handler_name.startswith("_"): raise ValueError("Malformed type in message (leading underscore)") return handler_name注意两个细节:消息必须携带type键,否则直接抛ValueError;类型名不允许以_开头,防止用户方法被下划线前缀的方法意外遮蔽。随后dispatch通过getattr(self, get_handler_name(message), None)找到处理器并调用;若找不到对应方法,会抛出ValueError("No handler for message type %s")(见 channels/consumer.py)。
你可能会问:我们怎么知道会收到哪些事件类型、里面有什么字段(比如websocket.receive里有text键)?答案是按 ASGI WebSocket 规范来设计——该规范定义了 WebSocket 如何被呈现,详见 docs/asgi.rst;再用一个检查scope["type"] == "websocket"的路由器保护这个应用,详见 docs/topics/routing.rst。
除此之外,最基本的 API 就是self.send(event):它把事件发回给客户端或协议服务器,具体语义由协议定义。对照 WebSocket 协议,上面发送的 dict 就是向客户端发一个文本帧。
2.3 AsyncConsumer 版本:一切皆协程
AsyncConsumer布局几乎一样,只是所有处理方法都必须是协程,self.send也是协程:
from channels.consumer import AsyncConsumer class EchoConsumer(AsyncConsumer): async def websocket_connect(self, event): await self.send({ "type": "websocket.accept", }) async def websocket_receive(self, event): await self.send({ "type": "websocket.send", "text": event["text"], })2.4 何时用 Sync、何时用 Async
主要看你在跟什么打交道:
- AsyncConsumer 内部若调用慢速同步函数,会阻塞整个事件循环,所以它只在你同时也在调用异步代码时才有价值(例如用
HTTPX并行抓取 20 个页面)。 - 只要涉及 Django ORM 或其他同步代码,就应使用 SyncConsumer,因为整个 consumer 会在一个线程里运行,避免 ORM 查询阻塞整个服务器。
官方推荐:默认优先写 SyncConsumer,只有在两种情况都满足时才用 AsyncConsumer——(1) 你明确知道当前工作能被异步处理改善(可并行的长任务);(2) 你只使用异步原生的库。
如果你确实想在 AsyncConsumer 里调用同步函数,可以借助asgiref.sync.sync_to_async——这正是 Channels 把 SyncConsumer 跑在线程池里所用的工具,它能把任何同步可调用对象包装成异步协程。
重要:若要在 AsyncConsumer(或任何异步代码)里调用 Django ORM,应当使用
database_sync_to_async适配器,或使用带a前缀的异步方法(如aget)。详见 docs/topics/databases.rst。
三、关闭消费者:StopConsumer 与资源清理
当连接被关闭——无论由你发起还是客户端发起——你通常会收到一个事件(如http.disconnect或websocket.disconnect),应用实例只有很短的时间去处理它。
完成断开后的清理工作后,需要抛出channels.exceptions.StopConsumer来干净地终止 ASGI 应用,让服务器回收它。如果不抛这个异常,服务器会等到应用关闭超时(Daphne 默认 10 秒),然后强制终止应用并发出警告。
下面这些通用消费者已经替你做了这件事,所以只有当你基于AsyncConsumer/SyncConsumer自己写消费者类时才需要手动处理。但如果你覆写了它们的__call__方法,或阻塞了它所调用的处理方法使其不返回,仍可能踩到这个坑——想深入了解可以阅读它们的源码。
另外,如果你启动了后台协程,务必在连接结束时把它们也关掉,否则会把协程泄漏给服务器。
从源码看,StopConsumer定义在 channels/exceptions.py,而AsyncConsumer.__call__用try/except StopConsumer捕获它来干净退出(channels/consumer.py)。事件循环本身由await_many_dispatch实现(channels/utils.py):它把传入的receive与 channel layer 的channel_receive都包装成任务,用asyncio.wait(..., return_when=FIRST_COMPLETED)监听,谁先完成就把结果交给dispatch,退出时统一取消所有任务。
四、Channel Layers:消费者之间的消息通道
Consumers 还支持 Channels 的channel layers(频道层),让消费者之间可以点对点发消息,或通过groups(组)进行广播。
消费者默认使用名为default的 channel layer;在子类化任何 Channels 提供的Consumer类时,可以设置channel_layer_alias属性来更换:
from channels.consumer import SyncConsumer class EchoConsumer(SyncConsumer): channel_layer_alias = "echo_alias"从 channels/consumer.py 可以看到,AsyncConsumer.__call__在初始化时通过get_channel_layer(self.channel_layer_alias)取得 channel layer,并调用new_channel()为当前消费者申请专属频道名self.channel_name;若 layer 为None(未配置),则退化为只从客户端receive读取消息。更多内容见 docs/topics/channel_layers.rst。
五、Scope:连接信息的载体
Consumers 在被调用时会收到连接的scope,里面包含大量你通常在 Django 视图request对象上能拿到的信息,在消费者方法内通过self.scope访问。
Scope 属于 ASGI 规范 的一部分,这里列出几个常用项:
| 键 | 含义 | 适用协议 |
|---|---|---|
scope["path"] | 请求的路径 | HTTP 与 WebSocket |
scope["headers"] | 请求的原始 name/value 头对 | HTTP 与 WebSocket |
scope["method"] | 请求使用的方法名 | HTTP |
如果启用了认证,你还能通过scope["user"]访问用户对象;URLRouter会把 URL 中捕获的分组放进scope["url_route"]。
概括地说,scope 是获取连接信息的地方,也是中间件放置属性供你访问的地方(这与 Django 中间件往request上挂东西的方式类似)。要查看连接 scope 的完整字段清单,需要查阅你所终止协议的基础 ASGI 规范,以及你使用的中间件/路由代码。
六、通用消费者(Generic Consumers)
上面展示的是适用于任意协议的基础布局。与 Django 的generic views类似,Channels 内置了通用消费者,把常见功能封装好,你无需重写——主要针对 HTTP 与 WebSocket 两类协议。
6.1 WebsocketConsumer(同步)
位于channels.generic.websocket.WebsocketConsumer,把啰嗦的裸 ASGI 消息收发封装成只处理文本帧和二进制帧的接口:
from channels.generic.websocket import WebsocketConsumer class MyConsumer(WebsocketConsumer): groups = ["broadcast"] def connect(self): # Called on connection. # To accept the connection call: self.accept() # Or accept the connection and specify a chosen subprotocol. # A list of subprotocols specified by the connecting client # will be available in self.scope['subprotocols'] self.accept("subprotocol") # To reject the connection, call: self.close() def receive(self, text_data=None, bytes_data=None): # Called with either text_data or bytes_data for each frame # You can call: self.send(text_data="Hello world!") # Or, to send a binary frame: self.send(bytes_data="Hello world!") # Want to force-close the connection? Call: self.close() # Or add a custom WebSocket error code! self.close(code=4123) def disconnect(self, close_code): # Called when the socket closes连接接受与拒绝:你也可以在connect方法的任意位置抛出channels.exceptions.AcceptConnection或channels.exceptions.DenyConnection来接受或拒绝连接——如果你想要可复用的、不依赖 mixin 的认证或限流代码,这种方式非常有用。这两个异常定义于 channels/exceptions.py,源码 channels/generic/websocket.py 会在websocket_connect中捕获它们并分别调用self.accept()/self.close()。
groups 自动加入/退出:WebsocketConsumer的频道会在连接时自动加入、断开时自动移出groups类属性列出的所有组。groups必须是可迭代对象,且必须配置支持组的 channel layer 作为后端(channels.layers.InMemoryChannelLayer与channels_redis.core.RedisChannelLayer都支持组)。如果未配置 channel layer,或 channel layer 不支持组,连接一个groups非空的WebsocketConsumer会抛出channels.exceptions.InvalidChannelLayerError。这一行为在源码 channels/generic/websocket.py 的websocket_disconnect(以及连接时的websocket_connect)中可以看到:group_add/group_discard遇到AttributeError时就会转抛InvalidChannelLayerError。组的详细说明见 docs/topics/channel_layers.rst。
6.2 AsyncWebsocketConsumer(异步)
位于channels.generic.websocket.AsyncWebsocketConsumer,与WebsocketConsumer的方法和签名完全相同,但一切都是异步的,你需要编写的方法也必须是协程:
from channels.generic.websocket import AsyncWebsocketConsumer class MyConsumer(AsyncWebsocketConsumer): groups = ["broadcast"] async def connect(self): # Called on connection. # To accept the connection call: await self.accept() # Or accept the connection and specify a chosen subprotocol. # A list of subprotocols specified by the connecting client # will be available in self.scope['subprotocols'] await self.accept("subprotocol") # To reject the connection, call: await self.close() async def receive(self, text_data=None, bytes_data=None): # Called with either text_data or bytes_data for each frame # You can call: await self.send(text_data="Hello world!") # Or, to send a binary frame: await self.send(bytes_data="Hello world!") # Want to force-close the connection? Call: await self.close() # Or add a custom WebSocket error code! await self.close(code=4123) async def disconnect(self, close_code): # Called when the socket closes底层协议转换一览(来自 channels/generic/websocket.py 源码):
accept(subprotocol=None, headers=None):发送{"type": "websocket.accept", "subprotocol": subprotocol},可选携带headers(L189-L196);receive(text_data=None, bytes_data=None):websocket_receive会把消息拆成text或bytes后调用(L198-L206);send(text_data=None, bytes_data=None, close=False):只发文本或只发二进制,两者都传或都不传会抛ValueError;传close=True时发送后立即关闭(L214-L225);close(code=None, reason=None):发送{"type": "websocket.close"},可带自定义code与reason(L227-L236)。
6.3 JsonWebsocketConsumer(同步,自动 JSON)
位于channels.generic.websocket.JsonWebsocketConsumer,工作方式与WebsocketConsumer相同,区别在于它会对 WebSocket 文本帧自动做 JSON 编解码。仅有的 API 差异:
- 你的
receive_json方法必须接收一个参数content,即解码后的 JSON 对象; self.send_json只接收一个参数content,会为你编码成 JSON。
若要定制 JSON 编解码,可覆写encode_json与decode_json两个类方法。默认实现就是json.dumps/json.loads(见 channels/generic/websocket.py)。
6.4 AsyncJsonWebsocketConsumer(异步,自动 JSON)
JsonWebsocketConsumer的异步版本,位于channels.generic.websocket.AsyncJsonWebsocketConsumer。注意:连encode_json和decode_json也是异步函数(channels/generic/websocket.py)。
6.5 AsyncHttpConsumer(异步 HTTP)
位于channels.generic.http.AsyncHttpConsumer,提供实现 HTTP 端点所需的基础原语:
from channels.generic.http import AsyncHttpConsumer class BasicHttpConsumer(AsyncHttpConsumer): async def handle(self, body): await asyncio.sleep(10) await self.send_response(200, b"Your response bytes", headers=[ (b"Content-Type", b"text/plain"), ])你需要自己实现handle方法。该方法收到整个请求体(单个 bytestring)。Headers 可以传元组列表,也可以传字典。响应体必须是 bytestring。
你还可以实现disconnect方法,在断开时执行清理(例如关闭你启动的协程)。它即使在非正常断开时也会运行,所以不要指望此时handle已经干净地执行完。
更底层的原语:如果需要对响应做更多控制(例如实现长轮询 long polling),应改用self.send_headers和self.send_body。下面这个示例已经用到了 channel layers(稍后详解):
import json from channels.generic.http import AsyncHttpConsumer class LongPollConsumer(AsyncHttpConsumer): async def handle(self, body): await self.send_headers(headers=[ (b"Content-Type", b"application/json"), ]) # Headers are only sent after the first body event. # Set "more_body" to tell the interface server to not # finish the response yet: await self.send_body(b"", more_body=True) async def chat_message(self, event): # Send JSON and finish the response: await self.send_body(json.dumps(event).encode("utf-8"))关键机制(见 channels/generic/http.py):
send_headers(*, status=200, headers=None):发送{"type": "http.response.start", "status": status, "headers": headers}。注意 ASGI 规范要求协议服务器只有在你第一次调用send_body之后才开始向客户端发送响应;send_body(body, *, more_body=False):发送{"type": "http.response.body", "body": body, "more_body": more_body}。more_body=True表示后面还有内容、响应不结束;默认行为会结束响应,此后该 channel 上的消息将被忽略;send_response(status, body, **kwargs):send_headers+send_body的薄封装,只能调用一次。
这些原语同样可以用来实现Server-Sent Events(SSE):
from datetime import datetime from channels.generic.http import AsyncHttpConsumer class ServerSentEventsConsumer(AsyncHttpConsumer): async def handle(self, body): await self.send_headers(headers=[ (b"Cache-Control", b"no-cache"), (b"Content-Type", b"text/event-stream"), (b"Transfer-Encoding", b"chunked"), ]) while True: payload = "data: %s\n\n" % datetime.now().isoformat() await self.send_body(payload.encode("utf-8"), more_body=True) await asyncio.sleep(1)HTTP 消费者的事件流(channels/generic/http.py)由http_request(把分片 body 拼接完成后交给handle,最后统一调用disconnect并抛StopConsumer)与http_disconnect(做清理后抛StopConsumer)组成。
七、把消费者变成 ASGI 应用:as_asgi()
AsyncConsumer提供了类方法as_asgi(**initkwargs)(channels/consumer.py),返回一个 ASGI v3 单可调用对象:每个 scope 实例化一个 consumer 实例,作用类似于 Django 的as_view()。initkwargs会传给 consumer 构造函数。因此你可以在路由中这样使用:
# 例如与 URLRouter/ProtocolTypeRouter 配合(详见 docs/topics/routing.rst) application = URLRouter([ path("ws/chat/", MyConsumer.as_asgi()), ])八、用仓库测试验证你的消费者
仓库的测试代码是学习通用消费者行为的最佳参考。例如 tests/test_generic_websocket.py 中的test_websocket_consumer,用channels.testing.WebsocketCommunicator模拟客户端,完整走了一遍「连接 → 发文本 → 发二进制 → 断开」的流程:
@pytest.mark.django_db @pytest.mark.asyncio async def test_websocket_consumer(): results = {} class TestConsumer(WebsocketConsumer): def connect(self): results["connected"] = True self.accept() def receive(self, text_data=None, bytes_data=None): results["received"] = (text_data, bytes_data) self.send(text_data=text_data, bytes_data=bytes_data) def disconnect(self, code): results["disconnected"] = code app = TestConsumer() communicator = WebsocketCommunicator(app, "/testws/") connected, _ = await communicator.connect() assert connected assert "connected" in results await communicator.send_to(text_data="hello") response = await communicator.receive_from() assert response == "hello" assert results["received"] == ("hello", None) await communicator.send_to(bytes_data=b"w\0\0\0") response = await communicator.receive_from() assert response == b"w\0\0\0" assert results["received"] == (None, b"w\0\0\0") await communicator.disconnect() assert "disconnected" in resultsWebsocketCommunicator位于 channels/testing/websocket.py,它会自动为你构造websocket类型的 scope(含path、query_string、headers、subprotocols),并提供connect()(返回(是否接受, 子协议或关闭码))、send_to/send_json_to、receive_from/receive_json_from、disconnect()等快捷方法。异步消费者则可以用channels.testing.websocket.WebsocketCommunicator配合AsyncWebsocketConsumer做同样的测试。HTTP 端点的测试可参考 tests/test_generic_http.py 与 channels/testing/http.py。
九、参考与延伸阅读
- 事件分发与生命周期实现:channels/consumer.py
- 通用消费者实现:channels/generic/websocket.py、channels/generic/http.py
- 异常定义(
StopConsumer、AcceptConnection、DenyConnection、InvalidChannelLayerError):channels/exceptions.py - 事件循环
await_many_dispatch:channels/utils.py - 消费者测试工具:channels/testing/websocket.py、channels/testing/http.py
- 相关主题文档:路由、Channel Layers、数据库、认证、ASGI 规范
- 本文对应仓库版本:Channels 4.2.0(见 channels/init.py)
- 后端
- WebSocket
- 异步编程
【免费下载链接】channels
Developer-friendly asynchrony for Django
相关推荐
Django Channels消费者编写指南:从同步到异步的完整实践
Django Channels消费者编写指南:从同步到异步的完整实践 想要为你的Django项目添加实时通信功能吗?Django Channels消费者是构建W
后端WebSocket异步编程kafka-python 使用指南:从消费者到生产者的完整实践
kafka python 使用指南:从消费者到生产者的完整实践 概述 kafka python 是一个功能强大的 Python Kafka 客户端库,提供了与
后端消息队列Apache Pulsar 端到端消息加密实战:从密钥生成到生产者/消费者配置的完整指南
Apache Pulsar 端到端消息加密实战:从密钥生成到生产者/消费者配置的完整指南 导读 本文以 Apache Pulsar 官方 Cookbook 文档
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考