news 2026/9/27 10:40:11

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
不引 Redis 行不行?Litestar 的 PsycoPgChannelsBackend 把 PostgreSQL 变成广播中枢

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

【免费下载链接】litestarLight, flexible and extensible ASGI framework | Built to scale项目地址: https://gitcode.com/GitHub_Trending/li/litestar

Litestar 的 Channels 子系统里藏着一个"数据库即 broker"的后端:PsycoPgChannelsBackend。它把 PostgreSQL 原生的LISTEN/NOTIFY封装成标准ChannelsBackend,如果你的项目本来就在用 psycopg3,又多进程应用需要互相同步消息,一个 DSN 就能拿到跨进程事件广播,不用额外跑 Redis。这篇文章从启动到关闭、从故障到选型,把整个实现讲透。

它适合谁用?

先想一个具体场景:管理后台改了一条配置,希望所有在线的 WebSocket 客户端都收到"配置已更新"。消息不大、频率不高,但进程可能部署了好几份。

这时你有两条路:

  • 引入 Redis,搭一套 Pub/Sub;
  • 或者把已经在用的 PostgreSQL 直接当消息中转站。

第二条路就是PsycoPgChannelsBackend的定位。它和 Channels 架构里其他角色(抽象基类 ChannelsBackend、管理订阅与路由的ChannelsPlugin、包装单条事件流的Subscriber)的关系,见 docs/usage/channels.rst 里的术语表和流程图,概念本身不复杂,这里不再展开。

两个前提要确认:

  1. 你的 psycopg3 版本 ≥ 3.2.4(文档中 Backends 一节有明确版本要求);
  2. 你能接受没有历史回放——这点后文单独说,它是这个后端最大的限制。

三步接入 🔌

整个接入就是"建后端 → 挂插件 → 用插件",三步:

from litestar.channels import ChannelsPlugin from litestar.channels.backends.psycopg import PsycoPgChannelsBackend backend = PsycoPgChannelsBackend(pg_dsn="postgresql://user:pass@localhost:5432/mydb") channels = ChannelsPlugin( backend=backend, channels=["general", "notifications"], create_ws_route_handlers=True, )

三步拆开看:

  • 建后端:构造函数只收一个pg_dsn。你不用操心连接细节——源码里所有连接都是AsyncConnection.connect(self._pg_dsn, autocommit=True)建出来的,LISTEN/NOTIFY不需要事务包裹,autocommit 是刻意为之;
  • 挂插件:ChannelsPlugin会把自己注册为应用依赖(注入键名就是channels),并在应用启动/关闭时替你调用backend.on_startup()/on_shutdown()。生命周期你完全不用管;
  • 用插件:在 handler 里直接注入channels,调channels.publish(data, "general")发消息;create_ws_route_handlers=True会为每个声明的频道自动生成 WebSocket 路由。频道名不在channels列表里就会抛ChannelsException,想动态开频道就设arbitrary_channels_allowed=True。

消息发出去之后,到底经历了什么?

把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, autocommit=True) as conn: for channel in channels: await conn.execute( SQL("NOTIFY {channel}, {data}").format(channel=Identifier(channel), data=dec_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)有三个设计点:

  1. 差集计算:只对"请求的频道 - 已订阅的频道"执行LISTEN,重复订阅天然幂等。unsubscribe用交集,逻辑对称;
  2. 全程持锁:所有订阅变更都在self._listener_lock里串行执行,_subscribed_channels集合和监听连接不会被并发修改。单测test_subscription_mutations_are_serialized用假连接统计并发操作数,断言它永远不超过 1,就是为了钉死这个约束;
  3. 变更前先停监听:LISTEN/UNLISTEN和notifies()迭代共用同一条监听连接,所以改订阅前要先_stop_listener(),改完在finally里(除非正在关闭)重新_start_listener()。这保证了"订阅集合"和"连接上实际监听的状态"永远一致。

从启动到关闭,都发生了什么

启动(on_startup)做四件事:

  • 新建AsyncExitStack和事件队列(重置状态,允许复用);
  • 连出专用监听连接_listener_conn——长期驻留,只负责收通知,登记进 exit stack,异常路径下也能被统一关闭;
  • 置位标志、启动后台监听协程_listen();
  • 插件侧随后会拿channels列表调一次subscribe,把初始订阅挂上。

关闭(on_shutdown)反着来,且在锁内执行:

  1. 置_shutting_down = True(后续订阅变更的finally分支看到它就不会再重启监听);
  2. _stop_listener()停掉监听协程;
  3. 清空订阅集合;
  4. 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(timeout=0.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_history:

  • channels.subscribe(..., history=N);
  • channels.put_subscriber_history(...);
  • 生成 WebSocket 路由时设的ws_handler_send_history。

如果你的需求是"客户端连上来先补发最近 N 条",换后端就行:RedisChannelsStreamBackend(基于 Redis Streams)或MemoryChannelsBackend(history=20)都支持。同为 PostgreSQL 实现的AsyncPgChannelsBackend(asyncpg 驱动)也一样不支持历史,所以这不是 psycopg 驱动特有的缺陷,而是LISTEN/NOTIFY机制的天花板。

五个后端放一起看

后端消息介质历史回放一句话选型建议
MemoryChannelsBackend进程内存支持(history参数)单进程、测试、本地开发;速度最快
RedisChannelsStreamBackendRedis Streams支持要补发历史,Redis 阵营首选
RedisChannelsPubSubBackendRedis Pub/Sub不支持低延迟纯扇出,对延迟最敏感
PsycoPgChannelsBackendPostgreSQL(psycopg3)不支持已深度使用 psycopg3,不想加中间件
AsyncPgChannelsBackendPostgreSQL(asyncpg)不支持已深度使用 asyncpg;功能与上一行对齐,参数名不同

和 asyncpg 版对比,psycopg 版的构造差异只有参数名(pg_dsnvsdsn/make_connection),底层都是同一条LISTEN/NOTIFY语义。

决策前清单 ✅

下手前过一遍这五条,能省掉大部分返工:

  1. 版本:psycopg3 是否 ≥ 3.2.4?低了先升级再谈;
  2. 历史:产品上有没有"连上补发最近消息"的需求?有就选 Redis Streams 或带history的内存后端;
  3. 吞吐:广播是低频通知还是高吞吐推送?后者 Redis Pub/Sub 通常更合适,NOTIFY适合"数据库即 broker"的轻量场景;
  4. 部署形态:多实例是否都能拿到同一个 DSN?跨进程一致性完全由 PostgreSQL 保证,连不上库的进程收不到任何消息;
  5. 送达时机:有没有"必须确认已发布"的调用点?有的地方记得换成wait_published,publish只管入队。

满足这些,这个百来行的后端就能稳定地替你干"多进程广播"这件事:连接交给AsyncExitStack,订阅变更交给锁,故障交给事件队列显式上报,退订竞态靠二次校验兜底。

【免费下载链接】litestarLight, flexible and extensible ASGI framework | Built to scale项目地址: https://gitcode.com/GitHub_Trending/li/litestar

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/27 10:39:12

会聊天的机器人为什么还需要一颗STM32?从系统架构到工程实践

1. 一颗STM32在“会聊天的机器人”里到底扛了什么活很多人第一次看到“会聊天的机器人”这个词,脑子里浮现的画面大概是这样的:一个圆头圆脑的小家伙,能听懂你说话,能跟你插科打诨,甚至还能在你心情不好的时候讲个冷笑…

作者头像 李华
网站建设 2026/9/27 10:36:10

OpenHarmony I2C总线开发实战:协议机制、HDF驱动适配与排障指南

1. 从一根线说起:I2C 在 OpenHarmony 里到底扮演什么角色搞 OpenHarmony 设备开发的朋友,绕不开的一个话题就是外设接入。你拿到一块 RK3568 或者 Hi3861 的开发板,想把温湿度传感器、OLED 屏、EEPROM、触摸芯片这些外设接上去,第…

作者头像 李华
网站建设 2026/9/27 10:35:45

大模型基础概念

本质:LLM 是一个“概率预测机” 核心启示: LLM 实际上并不“知道”事实,它只是在模仿训练数据中词语出现的统计规律。这就是“幻觉”(Hallucination)的根源——它可能自信地输出了一个概率很高但逻辑错误的词。 关键参数&#x…

作者头像 李华
网站建设 2026/9/27 10:34:25

STM32+Air780E+OLED:按键触发中文短信发送终端实战

1. 项目缘起与整体方案拆解按键一按,短信发出,OLED屏幕上实时滚动着“发送中”“发送成功”的状态——这个场景听起来像是某个工业设备的报警通知模块,或者是一个远程数据采集终端的核心交互逻辑。我最近刚把一个类似的项目从零跑通&#xff…

作者头像 李华
网站建设 2026/9/27 10:33:10

VSCode离线配置ESP32开发环境:ESP-IDF多版本共存与新项目向导实战

1. 为什么这个教程值得你花30分钟认真读完VSCODE安装ESP32开发环境,表面看只是点几下鼠标、敲几行命令的事,但实际踩过的坑,足够让一个有C语言基础的工程师在头三天反复重启电脑、重装系统、怀疑人生。我带过6个应届生做物联网毕设&#xff0…

作者头像 李华