80 lines
2.0 KiB
Python
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())
|