不引入 Redis,跨进程事件广播也能做:Litestar Channels psycopg 后端完全指南
【免费下载链接】litestarLight, flexible and extensible ASGI framework | Built to scale项目地址: https://gitcode.com/GitHub_Trending/li/litestar
服务里已经有 PostgreSQL,难道为了跨进程广播消息还要再养一套 Redis?Litestar 的 Channels 子系统给出了一个更克制的选项:借助LISTEN/NOTIFY,让数据库本身充当消息 broker。本文以 psycopg 后端PsycoPgChannelsBackend为例,跟着一条消息走完从发布到消费的完整链路,把这套约 120 行实现的设计取舍与坑一次讲清。
无需 Redis 的跨进程广播怎么配
后端构造只吃一个参数pg_dsn:
from litestar.channels import ChannelsPlugin from litestar.channels.backends.psycopg import PsycoPgChannelsBackend channels = ChannelsPlugin( backend=PsycoPgChannelsBackend("postgresql://user:pass@localhost:5432/appdb"), channels=["general", "notifications"], create_ws_route_handlers=True, )两个约束先说清:驱动需要psycopg 3.2.4 及以上(内部依赖较新版本才有的异步notifies()接口,版本要求写在 docs/usage/channels.rst 的 Backends 一节);生命周期不用你管,ChannelsPlugin会在应用启动/关闭时替你调on_startup()/on_shutdown()。频道名需提前在channels=里声明,否则发布或订阅未声明频道会抛ChannelsException,除非开启arbitrary_channels_allowed=True。
跟一条消息走:从 NOTIFY 到订阅者
先给结论:发布走短连接,接收走长连接,中间用队列中转。理解了这三个分工,其余代码都是推论。
发布端:每个频道一条 NOTIFY
它解决的问题:NOTIFY的载荷只能是文本,而频道名是动态值,不能硬拼进 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) )data: bytes先decode("utf-8")成文本,频道名经Identifier转义,避免注入(源码见 litestar/channels/backends/psycopg.py)。取舍在于:发布不复用监听连接,而是每次新建短连接。好处是发消息不会挤占接收通道;代价是每次发布多一次建连,低频通知场景完全可接受。逐频道执行NOTIFY,也天然完成了 fanout——所有LISTEN该频道的连接,无论属于哪个进程,都会收到。
监听端:为什么专门留一条长连接
它解决的问题:psycopg3 的notifies()异步迭代器和LISTEN/UNLISTEN语句共享同一条连接,监听中途执行别的语句会互相干扰。
监听是一个常驻后台任务:
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")))_LISTEN_POLL_INTERVAL是 0.1 秒:每轮notifies()只跑这么久就返回一次,循环借此定期复查"是否该停"。取舍是失败不沉默——除CancelledError原样重抛外,任何异常(比如连接断了)都会被塞进事件队列,和正常事件走同一条通道,由消费者负责区分。停监听则"先礼后兵":置位_stop_listening后最多等_STOP_LISTENER_TIMEOUT(5 秒)让循环自行退出,超时才cancel()任务。
订阅端:差集、锁,以及为什么改之前要先停
它解决的问题:并发的subscribe/unsubscribe会同时改同一个集合和同一条连接,稍不注意订阅状态就乱了。
async with self._listener_lock: new = requested - self._subscribed_channels if not new: return await self._stop_listener() for channel in new: await self._listener_conn.execute(SQL("LISTEN {channel}").format(channel=Identifier(channel))) self._subscribed_channels.add(channel)三个取舍:只对"请求集合 − 已订阅集合"执行LISTEN,重复订阅天然幂等;整个变更过程被_listener_lock串行化;因为改订阅和收事件共用一条连接,变更期间必须先停掉监听任务,finally里若不在关闭中则重新启动。unsubscribe是镜像实现,用交集算出要UNLISTEN的目标。
消费端:退订后的残留事件为什么会被丢弃
它解决的问题:退订瞬间,队列里可能还躺着该频道的旧事件,直接放行就会把已退订频道的数据交给订阅者。
while True: event = await self._event_queue.get() if isinstance(event, Exception): raise event if event[0] in self._subscribed_channels: yield eventstream_events()是插件后台 worker 的消费入口,取出的每个事件做二次校验:是Exception实例就直接抛——监听故障就是这样抵达消费者的;频道已不在订阅集合里则静默丢弃。取舍是极简:队列把网络 IO 和业务消费速率彻底解耦,代价只是每个事件多一次集合成员检查。
边界与坑
按"现象 → 原因 → 规避"组织:
- 历史回放不可用:
get_history()直接抛NotImplementedError,NOTIFY是"发出即散"的机制,库里不存频道历史。需要 WebSocket 新连接时补发历史的话,改用RedisChannelsStreamBackend或内存后端(相关用法见 docs/examples/channels/put_history.py)。 - 驱动版本别降级:
notifies()的异步迭代要求 psycopg ≥ 3.2.4,锁依赖时注意别把psycopg包意外降到旧版。 publish不等于"已送达":插件层publish()是同步非阻塞的,消息先入内部队列,由后台 worker 异步写入后端;需要确认已发布时用wait_published(),它绕过内部队列直接调backend.publish(见 litestar/channels/plugin.py)。- 监听断了不会自愈:连接中断以异常对象入队,最终由
stream_events抛给消费者——故障可感知,但监听任务不会自动重连,恢复要靠上层处理或重启进程。 - 未声明频道直接拒绝:不开
arbitrary_channels_allowed时,任何未声明频道的发布/订阅都是ChannelsException,这是特性不是 bug。
五种 Channels 后端选型对照
| 后端 | 消息 broker | 能否回放历史 | 一句话点评 |
|---|---|---|---|
MemoryChannelsBackend | 进程内存 | 能(history参数) | 单进程最快,测试与本地开发首选 |
RedisChannelsPubSubBackend | Redis Pub/Sub | 不能 | 低延迟跨进程广播的默认项 |
RedisChannelsStreamBackend | Redis Streams | 能 | 跨进程且要历史/持久化 |
AsyncPgChannelsBackend | PostgreSQL(asyncpg) | 不能 | asyncpg 技术栈的项目 |
PsycoPgChannelsBackend | PostgreSQL(psycopg3) | 不能 | psycopg3 技术栈的项目 |
选型建议:
- 广播场景低频(通知、后台刷新)时,复用现有 PostgreSQL 最省:少一个要运维的进程,跨实例一致性由数据库保证;
- 吞吐和尾延迟是硬指标时优先 Redis 系——无历史需求选 Pub/Sub,有则选 Streams;
- 两个 PG 后端功能对齐,差异仅在驱动与构造参数(
pg_dsn对dsn/make_connection),按项目已有驱动选即可。
用假连接验证后端行为
不启动 PostgreSQL 也能验证这个后端的内部行为。tests/unit/test_channels/test_psycopg_backend.py 直接把假连接对象注入_listener_conn,覆盖三个关键行为:
test_subscription_mutations_are_serialized:用asyncio.gather并发发起两次subscribe,断言假连接上"活跃操作数"峰值恒为 1,证明锁串行化生效;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 里的postgres_psycopg_backendfixture 会拉起 Docker PostgreSQL;结合 tests/unit/test_channels/test_backends.py 可以看到各后端共享的统一行为契约是怎么测的。
收束
PsycoPgChannelsBackend是ChannelsBackend契约的"数据库即 broker"实现:短连接发布、长连接监听、队列中转、锁串行订阅变更、故障经队列显式上抛——每一处都对应一个具体的并发或一致性问题。它的边界同样清晰:无历史回放、监听不自动重连。若广播量上升或需要新连接补发历史,下一步是切到 Redis 系后端,并顺手把max_backlog/backlog_strategy的背压策略配上。
【免费下载链接】litestarLight, flexible and extensible ASGI framework | Built to scale项目地址: https://gitcode.com/GitHub_Trending/li/litestar
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考