You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.27 17:10:38