ARTICLE DETAIL

资讯详情

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

不引 Redis 行不行?Litestar 的 PsycoPgChannelsBackend 把 PostgreSQL 变成广播中枢

不引 Redis 行不行?Litestar 的 PsycoPgChannelsBackend 把 PostgreSQL 变成广播中枢 不引 Redis 行不行Litestar 的 PsycoPgChannelsBackend 把 PostgreSQL 变成广播中枢【免费下载链接】litestarLight, flexible and extensible ASGI framework | Built to scale项目地址: https://gitcode.com/GitHub_Trending/li/litestarLitestar 的 Channels 子系统里藏着一个数据库即 broker的后端PsycoPgChannelsBackend。它把 PostgreSQL 原生的LISTEN/NOTIFY封装成标准ChannelsBackend如果你的项目本来就在用 psycopg3又多进程应用需要互相同步消息一个 DSN 就能拿到跨进程事件广播不用额外跑 Redis。这篇文章从启动到关闭、从故障到选型把整个实现讲透。它适合谁用先想一个具体场景管理后台改了一条配置希望所有在线的 WebSocket 客户端都收到配置已更新。消息不大、频率不高但进程可能部署了好几份。这时你有两条路引入 Redis搭一套 Pub/Sub或者把已经在用的 PostgreSQL 直接当消息中转站。第二条路就是PsycoPgChannelsBackend的定位。它和 Channels 架构里其他角色抽象基类 ChannelsBackend、管理订阅与路由的ChannelsPlugin、包装单条事件流的Subscriber的关系见 docs/usage/channels.rst 里的术语表和流程图概念本身不复杂这里不再展开。两个前提要确认你的 psycopg3 版本 ≥ 3.2.4文档中 Backends 一节有明确版本要求你能接受没有历史回放——这点后文单独说它是这个后端最大的限制。三步接入 整个接入就是建后端 → 挂插件 → 用插件三步from litestar.channels import ChannelsPlugin from litestar.channels.backends.psycopg import PsycoPgChannelsBackend backend PsycoPgChannelsBackend(pg_dsnpostgresql://user:passlocalhost:5432/mydb) channels ChannelsPlugin( backendbackend, channels[general, notifications], create_ws_route_handlersTrue, )三步拆开看建后端构造函数只收一个pg_dsn。你不用操心连接细节——源码里所有连接都是AsyncConnection.connect(self._pg_dsn, autocommitTrue)建出来的LISTEN/NOTIFY不需要事务包裹autocommit 是刻意为之挂插件ChannelsPlugin会把自己注册为应用依赖注入键名就是channels并在应用启动/关闭时替你调用backend.on_startup()/on_shutdown()。生命周期你完全不用管用插件在 handler 里直接注入channels调channels.publish(data, general)发消息create_ws_route_handlersTrue会为每个声明的频道自动生成 WebSocket 路由。频道名不在channels列表里就会抛ChannelsException想动态开频道就设arbitrary_channels_allowedTrue。消息发出去之后到底经历了什么把LISTEN/NOTIFY想象成数据库的公告栏谁LISTEN了某个频道数据库就把NOTIFY的载荷原样贴给谁。同一台数据库上的所有进程、所有连接都共享这块公告栏——这就是跨进程广播的全部魔法来源。一条消息的完整旅程用箭头写出来就是你的代码 publish() → 插件内部队列同步非阻塞 → 后台 pub worker → backend.publish()开一条新短连接逐频道 NOTIFY → PostgreSQL 按 LISTEN 分发 → 专用监听连接上的 _listen 协程捞到通知 → 事件队列put_nowait → stream_events() 逐个 yield → 插件 sub worker 分发给各 Subscriber → 你的 WebSocket 客户端有两个地方值得留意发布用短连接每次publish都新开一条连接、发完即关不和接收通知的长连接抢资源。核心 SQL 长这样async with await AsyncConnection.connect(self._pg_dsn, autocommitTrue) as conn: for channel in channels: await conn.execute( SQL(NOTIFY {channel}, {data}).format(channelIdentifier(channel), datadec_data) )频道名走Identifier转义数据先decode(utf-8)变回文本NOTIFY载荷本质是字符串接收端再encode(utf-8)维持了 Channels 接口事件一律是 bytes的约定。插件层的publish()是同步的它只是把消息扔进内部队列就返回不保证已经落到数据库。需要确认送达时用await channels.wait_published(data, channels)它会绕过队列直接调backend.publish。一次订阅是怎么变成 LISTEN 语句的subscribe的实现litestar/channels/backends/psycopg.py 里的subscribe/unsubscribe有三个设计点差集计算只对请求的频道 - 已订阅的频道执行LISTEN重复订阅天然幂等。unsubscribe用交集逻辑对称全程持锁所有订阅变更都在self._listener_lock里串行执行_subscribed_channels集合和监听连接不会被并发修改。单测test_subscription_mutations_are_serialized用假连接统计并发操作数断言它永远不超过 1就是为了钉死这个约束变更前先停监听LISTEN/UNLISTEN和notifies()迭代共用同一条监听连接所以改订阅前要先_stop_listener()改完在finally里除非正在关闭重新_start_listener()。这保证了订阅集合和连接上实际监听的状态永远一致。从启动到关闭都发生了什么启动on_startup做四件事新建AsyncExitStack和事件队列重置状态允许复用连出专用监听连接_listener_conn——长期驻留只负责收通知登记进 exit stack异常路径下也能被统一关闭置位标志、启动后台监听协程_listen()插件侧随后会拿channels列表调一次subscribe把初始订阅挂上。关闭on_shutdown反着来且在锁内执行置_shutting_down True后续订阅变更的finally分支看到它就不会再重启监听_stop_listener()停掉监听协程清空订阅集合exit_stack.aclose()关闭监听连接。_stop_listener()本身是先礼后兵先设_stop_listening True让循环自己退出最多等 5 秒_STOP_LISTENER_TIMEOUT超时才cancel()强制取消并等它收尾。这个 5 秒窗口给了监听循环体面退出的机会避免每次订阅变更都靠取消任务硬切。监听器坏了会怎样监听协程_listen()的主循环长这样while not self._stop_listening: async for notify in self._listener_conn.notifies(timeout_LISTEN_POLL_INTERVAL): self._event_queue.put_nowait((notify.channel, notify.payload.encode(utf-8)))notifies(timeout0.1)每 0.1 秒_LISTEN_POLL_INTERVAL一轮既能及时捞到通知也给循环留下检查停止标志的机会。异常处理上分两类CancelledError原样重抛——被强制取消就取消不多事其他任何异常连接断了、数据库重启……被装进事件队列当作一条事件发给消费者。消费者stream_events()那边接住event await self._event_queue.get() if isinstance(event, Exception): raise event if event[0] in self._subscribed_channels: yield event第一段就是故障传播监听端死了异常顺着队列一路抛到插件的订阅 worker上层立刻感知到这条链路断了而不是静默丢消息。第二段处理另一个隐蔽竞态UNLISTEN发出去的瞬间队列里可能还躺着旧频道的通知。所以每条事件取出后都要再校验一次该频道现在还在订阅集合里吗不在就丢弃。单测test_stream_events_filters_queued_events_after_unsubscribe和test_stream_events_propagates_listener_failures分别钉死了这两个行为——整个测试文件用假连接对象伪造notifies()不依赖真实数据库是理解这套内部状态的最佳阅读入口tests/unit/test_channels/test_psycopg_backend.py。它做不了的一件事历史回放get_history()直接抛NotImplementedError。LISTEN/NOTIFY是纯广播语义数据库不会替你存过去贴过什么。受影响的入口有三个都会走到get_historychannels.subscribe(..., historyN)channels.put_subscriber_history(...)生成 WebSocket 路由时设的ws_handler_send_history。如果你的需求是客户端连上来先补发最近 N 条换后端就行RedisChannelsStreamBackend基于 Redis Streams或MemoryChannelsBackend(history20)都支持。同为 PostgreSQL 实现的AsyncPgChannelsBackendasyncpg 驱动也一样不支持历史所以这不是 psycopg 驱动特有的缺陷而是LISTEN/NOTIFY机制的天花板。五个后端放一起看后端消息介质历史回放一句话选型建议MemoryChannelsBackend进程内存支持history参数单进程、测试、本地开发速度最快RedisChannelsStreamBackendRedis Streams支持要补发历史Redis 阵营首选RedisChannelsPubSubBackendRedis Pub/Sub不支持低延迟纯扇出对延迟最敏感PsycoPgChannelsBackendPostgreSQLpsycopg3不支持已深度使用 psycopg3不想加中间件AsyncPgChannelsBackendPostgreSQLasyncpg不支持已深度使用 asyncpg功能与上一行对齐参数名不同和 asyncpg 版对比psycopg 版的构造差异只有参数名pg_dsnvsdsn/make_connection底层都是同一条LISTEN/NOTIFY语义。决策前清单 ✅下手前过一遍这五条能省掉大部分返工版本psycopg3 是否 ≥ 3.2.4低了先升级再谈历史产品上有没有连上补发最近消息的需求有就选 Redis Streams 或带history的内存后端吞吐广播是低频通知还是高吞吐推送后者 Redis Pub/Sub 通常更合适NOTIFY适合数据库即 broker的轻量场景部署形态多实例是否都能拿到同一个 DSN跨进程一致性完全由 PostgreSQL 保证连不上库的进程收不到任何消息送达时机有没有必须确认已发布的调用点有的地方记得换成wait_publishedpublish只管入队。满足这些这个百来行的后端就能稳定地替你干多进程广播这件事连接交给AsyncExitStack订阅变更交给锁故障交给事件队列显式上报退订竞态靠二次校验兜底。【免费下载链接】litestarLight, flexible and extensible ASGI framework | Built to scale项目地址: https://gitcode.com/GitHub_Trending/li/litestar创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表