ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Litestar Channels psycopg 后端实战指南:用 PostgreSQL LISTEN/NOTIFY 做零额外中间件事件广播

Litestar Channels psycopg 后端实战指南:用 PostgreSQL LISTEN/NOTIFY 做零额外中间件事件广播 Litestar Channels psycopg 后端实战指南用 PostgreSQL LISTEN/NOTIFY 做零额外中间件事件广播【免费下载链接】litestarLight, flexible and extensible ASGI framework | Built to scale项目地址: https://gitcode.com/GitHub_Trending/li/litestarLitestar 的 Channels 子系统负责事件流路由最常见的用途是把一条消息广播给多个进程里的 WebSocket 客户端PsycoPgChannelsBackend是其中基于 psycopg3 异步驱动的实现把已有的 PostgreSQL 直接当作消息 broker靠原生的LISTEN/NOTIFY完成跨进程投递。如果你的项目重度依赖 psycopg3又不想为低频广播单独引入一个 Redis这篇 Litestar Channels psycopg 后端接入与原理指南讲的就是它怎么配置、一条消息如何从NOTIFY走到 WebSocket、有哪些坑、以及和其他后端怎么对比。多实例 WebSocket 广播不接 broker 会卡在哪先说清楚不接外部 broker 时的症状MemoryChannelsBackend把所有数据存在进程内存里性能是所有后端里最高的但进程之间互不可见。应用一旦横向扩容到多进程、多实例A 进程发布的事件到不了 B 进程的连接上的客户端——广播直接失效。Channels 的解法是把消息存储与扇出抽象成ChannelsBackend定义在 litestar/channels/backends/base.py由它负责和 broker 通信发布、接收。所有共享同一 broker 的后端实例都能收到相同的消息跨进程通信就此成立。PsycoPgChannelsBackend是这个抽象在 psycopg3 上的实现broker 就是你现有的 PostgreSQL——跨进程一致性由数据库保证任何进程只要LISTEN了某个频道就能收到任意进程发出的NOTIFY。接入前的检查项只有一条psycopg3 版本必须 ≥ 3.2.4这是 docs/usage/channels.rst 在 Backends 一节的明确要求。3 分钟接入步骤dsn 与 ChannelsPlugin 怎么配接入配置只有三行核心关键词就一个pg_dsn。backend PsycoPgChannelsBackend( pg_dsnpostgresql://user:passwordlocalhost:5432/mydb, ) channels ChannelsPlugin( backendbackend, channels[general, notifications], create_ws_route_handlersTrue, )为什么这样写构造函数只收pg_dsn一个参数插件拿到channels声明后create_ws_route_handlersTrue会为每个频道自动生成一条 WebSocket 路由省去手写转发逻辑。生命周期你不用管——litestar/channels/plugin.py 里的ChannelsPlugin在应用启动时调用backend.on_startup()、关闭时调用backend.on_shutdown()启动阶段打开一条专用监听连接并拉起后台监听任务关闭阶段置停止标志、停掉监听、清空订阅集合并经AsyncExitStack统一释放连接异常路径下连接也不会泄漏。两个配置细节决定你能否少踩坑。第一后端内部所有连接都以autocommitTrue打开——LISTEN/NOTIFY这类指令不依赖事务上下文事务包裹反而没意义。第二频道名拼 SQL 时走 psycopg 的Identifier安全转义不会出现在文本字面量之外SQL 注入风险被工具链挡掉。另外记一条运行期行为插件层的publish()是同步非阻塞方法消息先进内部队列、由后台 worker 异步写入后端它返回时不保证已发布需要立即确认时用await channels.wait_published(data, channels)它会绕过队列直接调用backend.publish。 向未声明的频道发布或订阅会抛ChannelsException如果你的频道名要动态生成插件参数arbitrary_channels_allowedTrue可以允许临时建频道此时只生成一条用路径参数承载频道名的路由。一条消息的完整链路从 NOTIFY 到被消费发布端的实现很直白async def publish(self, data: bytes, channels: Iterable[str]) - None: dec_data data.decode(utf-8) async with await AsyncConnection[Any].connect(self._pg_dsn, autocommitTrue) as conn: for channel in channels: await conn.execute( SQL(NOTIFY {channel}, {data}).format(channelIdentifier(channel), datadec_data) )这段代码回答了两个为什么NOTIFY的载荷本质是文本所以入参bytes先decode(utf-8)接收端再encode还原守住事件一律为 bytes的接口约定每次发布新开一条短连接是因为监听连接被notifies()迭代长期占用复用它会和接收循环互相阻塞——逐频道NOTIFY天然就是 fanout所有LISTEN该频道的连接无论属于哪个进程都会收到。完整链路四步publish发出NOTIFY→ 监听任务收到通知并put_nowait进_event_queue→ 插件的订阅 worker 调stream_events()逐个消费 → 路由到对应频道的 Subscriber最终经 WebSocket 发给客户端。监听任务的节奏由两个常量控制_LISTEN_POLL_INTERVAL 0.1秒一轮notifies()每轮之间检查停止标志停止监听时先置_stop_listening True最多等_STOP_LISTENER_TIMEOUT 5.0秒让循环自己退出超时再cancel()任务即先礼后兵。错误传播走同一条队列监听循环捕获到异常后原样重抛CancelledError其他任何异常比如连接中断都作为对象放进事件队列stream_events从队列取到Exception实例时直接raise插件的订阅 worker 因此能感知故障并中止而不是静默丢消息。还有一处容易漏掉的竞态UNLISTEN生效前队列里可能还残留旧频道的事件。stream_events在取出事件后会再校验一次channel in self._subscribed_channels才 yield退订后的旧事件就此被丢弃——这是订阅集合与网络 IO 解耦后必须补的双重校验。不支持历史回放怎么办get_history 直接抛 NotImplementedError⚠️ 这个后端不支持频道历史回放get_history()直接抛NotImplementedError。落到使用上就是三条路全部堵死ChannelsPlugin.subscribe(..., historyN)会失败、put_subscriber_history()会失败、生成 WebSocket 路由时设置发送历史示例见 docs/examples/channels/put_history.py也会失败。所以你应该按需求换后端需要历史就用RedisChannelsStreamBackend基于 Redis Streams纯本地开发或单进程应用可以用MemoryChannelsBackend(history20)启用内存历史。同为 PostgreSQL 驱动的AsyncPgChannelsBackendlitestar/channels/backends/asyncpg.py在get_history上行为一致选型时别指望它补上这块能力。五个后端横向对比psycoPg 后端适合什么场景docs/usage/channels.rst 当前实现的后端对比如下后端Broker支持历史适用场景MemoryChannelsBackend进程内存支持history参数测试、本地开发、单进程应用RedisChannelsPubSubBackendRedis Pub/Sub否低延迟、无需历史的跨进程广播RedisChannelsStreamBackendRedis Streams支持需要历史与持久化的跨进程广播AsyncPgChannelsBackendPostgreSQLasyncpg 驱动否已深度使用 asyncpg 的项目PsycoPgChannelsBackendPostgreSQLpsycopg3 驱动否已深度使用 psycopg3 的项目选型建议按依赖现状走项目已经在用 psycopg3广播只是低频通知、管理端刷新这类需求就复用现有 PostgreSQL 基础设施减少一个中间件项目用的是 asyncpg则选AsyncPgChannelsBackend——两个后端功能对齐构造参数名不同pg_dsnvsdsn/make_connection。延迟和吞吐是硬指标时Redis Pub/Sub 通常是更优选择需要历史加持久化上 Redis Streams。不连数据库也能测假连接验证并发与错误处理这个后端的单元测试tests/unit/test_channels/test_psycopg_backend.py不连真实数据库而是注入假连接对象三个测试各验证一条关键契约test_subscription_mutations_are_serialized用asyncio.gather同时发起两次subscribe断言连接上活跃操作数峰值为 1——订阅变更在_listener_lock内串行执行不会并发写坏监听连接test_stream_events_propagates_listener_failures注入一个notifies()持续抛RuntimeError(listener failed)的连接断言该异常能从stream_events传播给消费者test_stream_events_filters_queued_events_after_unsubscribe队列里预置旧频道和新频道各一条事件断言退订后旧事件被过滤、新事件正常放行。假连接 队列断言的写法正好对应前文的内部状态机本地想验证并发与错误传播行为时照着抄即可配合 tests/unit/test_channels/conftest.py 和各后端共用的行为契约测试覆盖是完整的。边界结论项目已在用 psycopg3、广播是低频需求、不想新增中间件——PsycoPgChannelsBackend直接用确认 psycopg3 ≥ 3.2.4需要历史回放或更高吞吐——换RedisChannelsStreamBackend或RedisChannelsPubSubBackend不要在这个后端上硬做。【免费下载链接】litestarLight, flexible and extensible ASGI framework | Built to scale项目地址: https://gitcode.com/GitHub_Trending/li/litestar创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表