如何用Python异步迭代器实现观察者模式?现有方案是否最优?
基于asyncio实现类似Dart Broadcast Stream的变更流
问题描述
我正在尝试用Python的asyncio和异步迭代器实现观察者模式,目标是创建一个变更流(Change Stream),让任务可以添加变更,其他任务能以异步迭代器的形式订阅这些变更,希望实现类似Dart中broadcast streams的接口。以下是我目前的简化实现代码:
from asyncio import Condition class ChangeStream: def __init__(self): self._condition = Condition() self._change = None async def add_change(self, change): async with self._condition: self._change = change self._condition.notify_all() async def __aiter__(self): async with self._condition: while True: await self._condition.wait() yield self._change
需求补充:观察者需要能看到订阅之后的所有变更,看不到订阅前的变更。
请问这是否是实现观察者模式的最佳方案?有没有更高效、更易理解的实现方式?
当前实现的问题
- 变更丢失风险:如果多个变更连续触发,后面的变更会覆盖
_change变量,导致处理慢的观察者错过中间的变更。比如,观察者刚处理完一个变更,还没回到wait()状态时,新的变更被添加,此时_change被覆盖,观察者下次唤醒只能拿到最新的那个,中间的变更直接丢失。 - 订阅者冲突:所有订阅者共享同一个
_change,如果多个订阅者同时处理,可能因为变量被异步覆盖而拿到错误的内容。
更优实现方案
方案一:基于asyncio.Queue的广播流(推荐)
用队列给每个订阅者单独存变更,确保每个订阅者都能收到所有订阅后的变更,逻辑清晰且健壮。
import asyncio from typing import AsyncIterable, TypeVar T = TypeVar('T') class ChangeStream: def __init__(self): self._subscribers: list[asyncio.Queue[T]] = [] self._lock = asyncio.Lock() async def add_change(self, change: T): async with self._lock: # 给每个订阅者的队列推送变更 for queue in self._subscribers: await queue.put(change) def subscribe(self) -> AsyncIterable[T]: # 每个订阅者创建独立队列 queue = asyncio.Queue[T]() asyncio.create_task(self._add_subscriber(queue)) return self._stream_from_queue(queue) async def _add_subscriber(self, queue: asyncio.Queue[T]): async with self._lock: self._subscribers.append(queue) async def _stream_from_queue(self, queue: asyncio.Queue[T]) -> AsyncIterable[T]: try: while True: yield await queue.get() queue.task_done() finally: # 订阅者退出时自动清理队列 async with self._lock: self._subscribers.remove(queue)
优势
- 无变更丢失:每个订阅者的队列会缓存所有未处理的变更,即使处理速度慢也不会漏收。
- 订阅者独立:队列各自独立,不会出现变量覆盖的问题。
- 自动清理:订阅者停止迭代时会自动从订阅列表移除队列,避免内存泄漏。
- 逻辑直观:利用队列的天然特性实现消息分发,符合异步编程思维,容易理解和维护。
方案二:基于版本号+asyncio.Event的轻量实现
如果不需要缓存未处理的变更(只关心最新的,且确保订阅后不遗漏),可以用版本号跟踪解决覆盖问题:
import asyncio from typing import AsyncIterable, TypeVar T = TypeVar('T') class ChangeStream: def __init__(self): self._latest_change: T | None = None self._event = asyncio.Event() self._lock = asyncio.Lock() self._version = 0 async def add_change(self, change: T): async with self._lock: self._latest_change = change self._version += 1 self._event.set() self._event.clear() async def __aiter__(self) -> AsyncIterable[T]: current_version = self._version while True: await self._event.wait() async with self._lock: if self._version > current_version: current_version = self._version yield self._latest_change
这个方案通过版本号判断是否有新变更,每个订阅者跟踪自己的版本,避免了变更被覆盖导致的丢失问题,适合轻量场景。
总结
你的初始方案思路正确,但存在变更丢失和覆盖的问题。优先选择基于asyncio.Queue的实现,完全满足你“订阅后所有变更都能收到”的需求,且健壮性更强。如果追求轻量且只需要最新变更,版本号+事件的实现也可以。
内容的提问来源于stack exchange,提问作者odo
相关产品推荐
相关产品推荐

