ARTICLE DETAIL

资讯详情

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

不引入 Redis,跨进程事件广播也能做:Litestar Channels psycopg 后端完全指南

不引入 Redis,跨进程事件广播也能做:Litestar Channels psycopg 后端完全指南 不引入 Redis跨进程事件广播也能做Litestar Channels psycopg 后端完全指南【免费下载链接】litestarLight, flexible and extensible ASGI framework | Built to scale项目地址: https://gitcode.com/GitHub_Trending/li/litestar服务里已经有 PostgreSQL难道为了跨进程广播消息还要再养一套 RedisLitestar 的 Channels 子系统给出了一个更克制的选项借助LISTEN/NOTIFY让数据库本身充当消息 broker。本文以 psycopg 后端PsycoPgChannelsBackend为例跟着一条消息走完从发布到消费的完整链路把这套约 120 行实现的设计取舍与坑一次讲清。无需 Redis 的跨进程广播怎么配后端构造只吃一个参数pg_dsnfrom litestar.channels import ChannelsPlugin from litestar.channels.backends.psycopg import PsycoPgChannelsBackend channels ChannelsPlugin( backendPsycoPgChannelsBackend(postgresql://user:passlocalhost:5432/appdb), channels[general, notifications], create_ws_route_handlersTrue, )两个约束先说清驱动需要psycopg 3.2.4 及以上内部依赖较新版本才有的异步notifies()接口版本要求写在 docs/usage/channels.rst 的 Backends 一节生命周期不用你管ChannelsPlugin会在应用启动/关闭时替你调on_startup()/on_shutdown()。频道名需提前在channels里声明否则发布或订阅未声明频道会抛ChannelsException除非开启arbitrary_channels_allowedTrue。跟一条消息走从 NOTIFY 到订阅者先给结论发布走短连接接收走长连接中间用队列中转。理解了这三个分工其余代码都是推论。发布端每个频道一条 NOTIFY它解决的问题NOTIFY的载荷只能是文本而频道名是动态值不能硬拼进 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) )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_TIMEOUT5 秒让循环自行退出超时才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(channelIdentifier(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()直接抛NotImplementedErrorNOTIFY是发出即散的机制库里不存频道历史。需要 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参数单进程最快测试与本地开发首选RedisChannelsPubSubBackendRedis Pub/Sub不能低延迟跨进程广播的默认项RedisChannelsStreamBackendRedis Streams能跨进程且要历史/持久化AsyncPgChannelsBackendPostgreSQLasyncpg不能asyncpg 技术栈的项目PsycoPgChannelsBackendPostgreSQLpsycopg3不能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),仅供参考
返回列表