You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.14 12:31:00