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

如何在异步SQLAlchemy 2.x中使用Postgres NOTIFY/LISTEN或sqlalchemy事件?

异步实现PostgreSQL NOTIFY/LISTEN + SQLAlchemy方案

一、直接异步处理NOTIFY/LISTEN(基于SQLAlchemy AsyncEngine)

PostgreSQL的LISTEN/NOTIFY依赖单个数据库连接持续监听,在异步环境下可以通过SQLAlchemy的AsyncEngine单独维护监听连接,配合asyncio实现非阻塞监听。

实现代码:

import asyncio
from sqlalchemy.ext.asyncio import create_async_engine, AsyncConnection

async def listen_to_notifications(engine):
    # 建立独立的监听连接
    async with engine.connect() as conn:
        # 订阅目标频道
        await conn.execute("LISTEN user_updates;")
        # 获取底层asyncpg原生连接(SQLAlchemy异步引擎默认适配asyncpg)
        raw_conn = await conn.get_raw_connection()
        print("开始监听user_updates频道...")
        
        while True:
            # 异步等待通知,无通知时会挂起协程不阻塞事件循环
            notification = await raw_conn.wait_for_notify()
            if notification is None:
                break
            print(f"收到通知: 频道={notification.channel}, 内容={notification.payload}")

async def send_notification(engine, channel, payload):
    async with engine.begin() as conn:
        # 参数化执行避免SQL注入
        await conn.execute(f"NOTIFY {channel}, :payload;", {"payload": payload})

async def main():
    engine = create_async_engine("postgresql+asyncpg://user:pass@localhost/dbname")
    
    # 启动监听后台任务
    listen_task = asyncio.create_task(listen_to_notifications(engine))
    
    # 模拟发送通知
    await send_notification(engine, "user_updates", "user_123_updated")
    
    # 保持运行,实际场景可根据业务逻辑调整退出条件
    await listen_task

if __name__ == "__main__":
    asyncio.run(main())

二、结合SQLAlchemy ORM事件+异步通知

如果需要在ORM模型变更时自动触发NOTIFY,可以利用SQLAlchemy的ORM事件(如after_update),配合sqlalchemy.orm.attributes.get_history获取字段变更,再异步发送通知。

实现代码:

import asyncio
from sqlalchemy import Column, Integer, String
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
from sqlalchemy.orm import declarative_base, sessionmaker, attributes
from sqlalchemy import event

Base = declarative_base()

class User(Base):
    __tablename__ = "users"
    id = Column(Integer, primary_key=True)
    name = Column(String)
    email = Column(String)

async def on_user_updated(mapper, connection, target):
    # 获取字段变更历史
    name_history = attributes.get_history(target, "name")
    email_history = attributes.get_history(target, "email")
    
    # 仅当字段有变更时发送通知
    if name_history.has_changes() or email_history.has_changes():
        payload = f"user_id={target.id}, old_name={name_history.deleted[0] if name_history.deleted else '无'}, new_name={name_history.added[0] if name_history.added else '无'}"
        # 复用当前连接发送通知,自动纳入事务
        await connection.execute(f"NOTIFY user_updates, :payload;", {"payload": payload})

# 注册异步事件回调,必须指定asyncio=True以支持异步函数
event.listens_for(User, "after_update", asyncio=True)(on_user_updated)

async def main():
    engine = create_async_engine("postgresql+asyncpg://user:pass@localhost/dbname")
    # 初始化表结构
    async with engine.begin() as conn:
        await conn.run_sync(Base.metadata.create_all)
    
    AsyncSessionLocal = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
    
    # 启动监听任务(复用第一个方案中的listen_to_notifications函数)
    listen_task = asyncio.create_task(listen_to_notifications(engine))
    
    # 模拟更新用户触发事件
    async with AsyncSessionLocal() as session:
        user = await session.get(User, 1)
        if not user:
            user = User(id=1, name="OldName", email="old@example.com")
            session.add(user)
            await session.commit()
        
        user.name = "NewName"
        await session.commit()
    
    await asyncio.sleep(1)  # 等待通知被接收
    listen_task.cancel()
    await listen_task

if __name__ == "__main__":
    asyncio.run(main())

关键注意事项:

  • 监听连接生命周期:监听协程需保持运行,可作为后台常驻任务,避免随意关闭连接。
  • 异步事件回调:注册事件时必须指定asyncio=True,否则异步回调会被同步执行,阻塞事件循环。
  • 事务与NOTIFY:NOTIFY会在事务提交后才广播,after_update事件触发于事务提交前,该语句会被纳入当前事务,提交后自动生效。
  • 原生连接调用:通过get_raw_connection()可直接使用底层驱动(如asyncpg)的原生API处理通知,更高效。

内容的提问来源于stack exchange,提问作者Vladimir Alinsky

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 20:23:15