如何在Django/FastAPI中实现非阻塞式Redis长轮询?
在Django/FastAPI中实现非阻塞长轮询(结合Redis)
要搞非阻塞的长轮询,核心是靠异步IO处理Redis监听和HTTP请求,避免单个长连接占用线程/进程;同时用Redis的发布订阅(Pub/Sub)机制,数据库变更时主动给服务端推送事件,再由服务端响应客户端的长轮询请求。下面分框架给出具体实现方案:
FastAPI 实现(原生异步,推荐)
FastAPI基于Starlette,天生支持异步,配合异步Redis客户端可直接实现非阻塞长轮询,性能拉满。
安装依赖
pip install fastapi uvicorn redis
核心代码
from fastapi import FastAPI, HTTPException from redis.asyncio import Redis import asyncio app = FastAPI() # 初始化异步Redis客户端,默认使用连接池提升效率 redis = Redis(host="localhost", port=6379, db=0, decode_responses=True) # 定义Redis频道,数据库变更时往该频道推送消息 UPDATE_CHANNEL = "db_updates:user_data" @app.get("/long-poll") async def long_poll(client_id: str): if not client_id: raise HTTPException(status_code=400, detail="client_id 必填") pubsub = redis.pubsub() await pubsub.subscribe(UPDATE_CHANNEL) try: # 设置30秒超时,避免客户端永久挂起 timeout = 30 end_time = asyncio.get_event_loop().time() + timeout while asyncio.get_event_loop().time() < end_time: # 非阻塞等待消息,1秒超时避免卡住事件循环 msg = await pubsub.get_message(ignore_subscribe_messages=True, timeout=1) if msg: # 返回更新后的业务数据,msg['data']为数据库变更时推送的内容 return {"status": "updated", "data": msg["data"], "client_id": client_id} # 让出事件循环,给其他请求腾资源 await asyncio.sleep(0.1) # 超时未收到更新,返回标识让客户端自动重连 return {"status": "no_update", "client_id": client_id} finally: # 用完必须取消订阅、关闭连接,避免资源泄漏 await pubsub.unsubscribe(UPDATE_CHANNEL) await pubsub.close() # 模拟数据库变更触发消息(实际业务可挂载到ORM的保存钩子) @app.post("/trigger-update") async def trigger_update(data: str): await redis.publish(UPDATE_CHANNEL, data) return {"status": "success", "msg": "更新事件已推送"}
运行方式
用Uvicorn启动并开启多进程扛并发:
uvicorn main:app --workers 4 --host 0.0.0.0 --port 8000
Django 实现(异步视图支持)
Django 3.1+支持异步视图,配合ASGI服务器(Daphne/Uvicorn)可实现非阻塞长轮询。
安装依赖
pip install django redis
配置ASGI
在项目settings.py中指定ASGI应用:
ASGI_APPLICATION = "你的项目名.asgi.application"
异步视图代码
from django.http import JsonResponse from django.views import View from redis.asyncio import Redis import asyncio import json redis = Redis(host="localhost", port=6379, db=0, decode_responses=True) UPDATE_CHANNEL = "db_updates:user_data" class LongPollView(View): async def get(self, request): client_id = request.GET.get("client_id") if not client_id: return JsonResponse({"error": "client_id 必填"}, status=400) pubsub = redis.pubsub() await pubsub.subscribe(UPDATE_CHANNEL) try: timeout = 30 end_time = asyncio.get_event_loop().time() + timeout while asyncio.get_event_loop().time() < end_time: msg = await pubsub.get_message(ignore_subscribe_messages=True, timeout=1) if msg: return JsonResponse({ "status": "updated", "data": msg["data"], "client_id": client_id }) await asyncio.sleep(0.1) return JsonResponse({"status": "no_update", "client_id": client_id}) finally: await pubsub.unsubscribe(UPDATE_CHANNEL) await pubsub.close() # 数据库变更触发示例(使用Django信号) from django.db.models.signals import post_save from django.dispatch import receiver from .models import 你的业务模型 @receiver(post_save, sender=你的业务模型) async def send_db_update_signal(sender, instance, created, **kwargs): # 将更新后的模型数据转为JSON推送至Redis频道 update_data = json.dumps({ "id": instance.id, "name": instance.name, "updated_at": str(instance.updated_at) }) await redis.publish(UPDATE_CHANNEL, update_data)
运行方式
用Daphne启动ASGI服务:
daphne 你的项目名.asgi:application
关键注意事项
- Redis连接池:代码中Redis客户端默认使用连接池,不要手动创建单连接,否则高并发下会出现资源耗尽问题
- 超时控制:长轮询必须设置合理超时(30-60秒),客户端超时后自动重连,避免死连接
- 资源清理:每个请求的Pub/Sub对象必须关闭,否则Redis会堆积大量无效订阅
- 并发扩容:如果并发量极大,只需多部署几个服务实例,Redis的Pub/Sub会自动将消息广播给所有实例
内容的提问来源于stack exchange,提问作者Shukurullox Komiljonov
相关产品推荐
相关产品推荐

