如何在异步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
相关产品推荐
相关产品推荐

