Files
2026-08-02 14:07:04 +08:00

80 lines
2.0 KiB
Python

"""注释"""
# consumers.py
# 推荐解决方案:事件驱动架构
# 修改实现为事件驱动模式:
import asyncio
import json
import time
from channels.generic.http import AsyncHttpConsumer
# 全局消息队列 (生产环境用Redis等分布式队列)
MESSAGE_QUEUE = asyncio.Queue()
class SSEConsumer(AsyncHttpConsumer):
async def handle(self, body):
# 设置SSE响应头
await self.send_headers(
headers=[
(b"Content-Type", b"text/event-stream"),
(b"Cache-Control", b"no-cache"),
(b"Connection", b"keep-alive"),
]
)
# 事件循环监听
while True:
# 异步等待新消息(挂起状态)
data = await MESSAGE_QUEUE.get()
# 推送消息给客户端
await self.send_body(f"event: message\ndata: {data}\n\n".encode())
async def disconnect(self):
# 客户端断开时清理资源
print(f"Client {self.scope['client']} disconnected")
# 其他Django视图或后台任务中
# 生产者触发更新示例
async def trigger_update():
"""模拟外部事件触发更新"""
# 创建新消息
message = {"time": time.time(), "event": "DataUpdated"}
# 将消息推送到全局队列
await MESSAGE_QUEUE.put(json.dumps(message))
# 测试代码
async def xiaozizi():
"""从队列中获取消息并打印"""
print("等待消息...")
d = await MESSAGE_QUEUE.get()
print(f"收到消息: {d}")
return d
# 在脚本末尾添加(如果需要测试)
async def main():
"""主测试函数"""
print("开始测试...")
# 启动一个任务来监听消息
listen_task = asyncio.create_task(xiaozizi())
# 稍等一下再发送消息,确保监听任务已经启动
await asyncio.sleep(0.1)
# 发送测试消息
# await trigger_update()
# 等待监听任务完成
result = await listen_task
print(f"测试完成,结果: {result}")
if __name__ == "__main__":
asyncio.run(main())