Django中Asyncio监听PostgreSQL Notify阻塞HTTP/WebSocket请求问题
解决Django中PostgreSQL Notify监听阻塞主线程的问题
你的代码核心问题在于:在Django启动阶段(__init__.py中)调用loop.run_until_complete(db_listen())会完全阻塞主线程。run_until_complete会一直等待传入的协程执行完毕,而db_listen是无限循环的协程,直接导致Django无法处理任何HTTP或WebSocket请求。
下面是几种可行的解决方案,按场景推荐:
方案一:用独立线程运行数据库监听器(通用方案)
将数据库监听逻辑放到独立的守护线程中,让主线程专注处理Django的请求。同时注意psycopg2连接不是线程安全的,要给每个线程单独创建连接。
修改后的代码:
import threading import asyncio import psycopg2 from psycopg2.extensions import ISOLATION_LEVEL_AUTOCOMMIT # 替换成你自己的连接获取方法和WebSocket消费者 from your_app.utils import get_db_conn from your_app.consumers import NotificationConsumer # 异步队列,用于传递数据库通知 notify_queue = asyncio.Queue() async def process_notifications(): """持续从队列取通知并推送到WebSocket""" print('Started processing DB notifications ...') while True: notify = await notify_queue.get() await NotificationConsumer.send_data(notify.payload) print(f"Notification pushed: {notify.payload}") notify_queue.task_done() def handle_db_notify(connection, queue): """数据库通知回调:读取通知并放入异步队列""" def inner(): connection.poll() # 处理所有待处理的通知 while connection.notifies: notify = connection.notifies.pop(0) queue.put_nowait(notify) return inner def run_listener_thread(): """在独立线程中启动异步监听循环""" loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) # 为当前线程创建独立的数据库连接 conn = get_db_conn() conn.set_isolation_level(ISOLATION_LEVEL_AUTOCOMMIT) with conn.cursor() as cur: cur.execute("LISTEN new_item_added;") # 注册连接的读取回调 loop.add_reader(conn, handle_db_notify(conn, notify_queue)) # 启动通知处理协程 loop.run_until_complete(process_notifications()) def initiate_db_listener(): """启动监听线程(在app的__init__.py中调用)""" listener_thread = threading.Thread(target=run_listener_thread, daemon=True) listener_thread.start()
关键注意点:
- 设置
daemon=True,确保Django主进程退出时线程自动终止 - 每个线程必须使用独立的数据库连接,psycopg2连接不支持多线程共享
notify_queue是asyncio.Queue,只能在事件循环线程中调用put_nowait
方案二:基于Django Channels的全局监听(WebSocket场景推荐)
如果你的WebSocket是用Django Channels实现的,推荐用Channels的事件循环来处理数据库监听,避免线程管理的麻烦,同时可以方便地广播通知给所有WebSocket客户端。
1. 配置全局监听任务(在apps.py中)
from django.apps import AppConfig import asyncio from channels.db import database_sync_to_async import psycopg2 from psycopg2.extensions import ISOLATION_LEVEL_AUTOCOMMIT from channels.layers import get_channel_layer channel_layer = get_channel_layer() class YourAppConfig(AppConfig): default_auto_field = 'django.db.models.BigAutoField' name = 'your_app' def ready(self): # 启动全局数据库监听任务 asyncio.create_task(self.global_db_listener()) async def global_db_listener(self): """全局监听PostgreSQL通知并广播给WebSocket客户端""" @database_sync_to_async def create_listen_conn(): conn = psycopg2.connect( dbname="your_db", user="your_user", password="your_pass", host="your_host" ) conn.set_isolation_level(ISOLATION_LEVEL_AUTOCOMMIT) with conn.cursor() as cur: cur.execute("LISTEN new_item_added;") return conn # 创建监听连接 conn = await create_listen_conn() loop = asyncio.get_event_loop() def broadcast_notify(): conn.poll() while conn.notifies: notify = conn.notifies.pop(0) # 广播到名为"notifications"的Channel组 asyncio.run_coroutine_threadsafe( channel_layer.group_send( "notifications", { "type": "send_notification", "payload": notify.payload } ), loop ) # 注册连接读取回调 loop.add_reader(conn, broadcast_notify) try: # 保持协程存活 while True: await asyncio.sleep(3600) finally: loop.remove_reader(conn) conn.close()
2. 修改WebSocket消费者接收广播
from channels.generic.websocket import AsyncWebsocketConsumer class NotificationConsumer(AsyncWebsocketConsumer): async def connect(self): await self.accept() # 加入广播组 await self.channel_layer.group_add("notifications", self.channel_name) async def disconnect(self, close_code): # 退出广播组 await self.channel_layer.group_discard("notifications", self.channel_name) # 对应group_send中的"type"字段 async def send_notification(self, event): payload = event["payload"] await self.send(text_data=payload)
方案三:用Celery异步任务监听(已有Celery环境推荐)
如果你的项目已经在用Celery,可以将数据库监听逻辑写成Celery任务,由Celery Worker独立运行,完全不影响Django主线程。
1. 编写Celery监听任务
from celery import shared_task import time import psycopg2 from psycopg2.extensions import ISOLATION_LEVEL_AUTOCOMMIT from channels.layers import get_channel_layer from asgiref.sync import async_to_sync @shared_task(bind=True, autoretry_for=(Exception,), retry_backoff=3) def listen_db_notifications(self): conn = psycopg2.connect( dbname="your_db", user="your_user", password="your_pass", host="your_host" ) conn.set_isolation_level(ISOLATION_LEVEL_AUTOCOMMIT) with conn.cursor() as cur: cur.execute("LISTEN new_item_added;") try: while True: conn.poll() while conn.notifies: notify = conn.notifies.pop(0) # 通过Channel Layer广播通知 channel_layer = get_channel_layer() async_to_sync(channel_layer.group_send)( "notifications", { "type": "send_notification", "payload": notify.payload } ) # 避免CPU空转 time.sleep(1) except Exception as e: # 异常自动重试 self.retry(exc=e) finally: conn.close()
2. 在Django启动时触发任务(apps.py)
from django.apps import AppConfig import os class YourAppConfig(AppConfig): default_auto_field = 'django.db.models.BigAutoField' name = 'your_app' def ready(self): # 仅在主进程启动时触发,避免多Worker重复启动 if os.environ.get("RUN_MAIN") == "true": from your_app.tasks import listen_db_notifications listen_db_notifications.delay()
内容的提问来源于stack exchange,提问作者Mohit gupta
相关产品推荐
相关产品推荐

