怎样实现事件多观察者通知脚本

wen 实用脚本 30

本文目录导读:

怎样实现事件多观察者通知脚本

  1. 基础发布订阅模式
  2. 带优先级的通知系统
  3. 异步事件通知系统
  4. 通用解决方案 - 装饰器模式
  5. 实际应用示例
  6. 使用建议

我来帮你实现一个事件多观察者通知系统,这里提供几种不同场景的实现方案:

基础发布订阅模式

Python 实现

from typing import Any, Callable, Dict, List
from collections import defaultdict
class EventEmitter:
    """事件发射器 - 基础实现"""
    def __init__(self):
        self._observers: Dict[str, List[Callable]] = defaultdict(list)
    def on(self, event: str, callback: Callable) -> None:
        """订阅事件"""
        self._observers[event].append(callback)
        print(f"[订阅] 事件 '{event}' 添加了观察者: {callback.__name__}")
    def off(self, event: str, callback: Callable = None) -> None:
        """取消订阅"""
        if callback:
            if callback in self._observers[event]:
                self._observers[event].remove(callback)
                print(f"[取消订阅] 移除了事件 '{event}' 的观察者: {callback.__name__}")
        else:
            # 移除事件所有观察者
            self._observers[event].clear()
            print(f"[取消订阅] 清除了事件 '{event}' 的所有观察者")
    def emit(self, event: str, *args: Any, **kwargs: Any) -> None:
        """触发事件"""
        if event not in self._observers:
            print(f"[通知] 事件 '{event}' 没有观察者")
            return
        print(f"[通知] 触发事件 '{event}',通知 {len(self._observers[event])} 个观察者")
        for callback in self._observers[event]:
            try:
                callback(*args, **kwargs)
            except Exception as e:
                print(f"[错误] 观察者回调执行失败: {e}")
    def once(self, event: str, callback: Callable) -> None:
        """一次性订阅 - 只执行一次"""
        def wrapper(*args, **kwargs):
            callback(*args, **kwargs)
            self.off(event, wrapper)
        self.on(event, wrapper)
# 使用示例
if __name__ == "__main__":
    # 创建事件发射器
    emitter = EventEmitter()
    # 定义观察者
    def logger(data):
        print(f"[日志观察者] 收到数据: {data}")
    def email_sender(data):
        print(f"[邮件观察者] 发送邮件通知: {data}")
    def data_analyzer(data):
        print(f"[分析观察者] 分析数据: {data}")
    # 订阅事件
    emitter.on("data_updated", logger)
    emitter.on("data_updated", email_sender)
    emitter.on("data_updated", data_analyzer)
    # 一次性观察者
    emitter.once("data_updated", lambda data: print(f"[一次性] 首次处理: {data}"))
    # 触发事件
    print("\n=== 第一次触发 ===")
    emitter.emit("data_updated", {"user": "Alice", "action": "login"})
    print("\n=== 第二次触发 (一次性观察者已被移除) ===")
    emitter.emit("data_updated", {"user": "Bob", "action": "logout"})
    # 取消订阅
    print("\n=== 取消邮件观察者 ===")
    emitter.off("data_updated", email_sender)
    emitter.emit("data_updated", {"user": "Charlie", "action": "update"})

带优先级的通知系统

import heapq
from dataclasses import dataclass
from typing import Any, Callable
@dataclass(order=True)
class PriorityObserver:
    """带优先级的观察者"""
    priority: int
    callback: Callable = None
    def __post_init__(self):
        # 确保回调函数可比较
        self.callback_name = self.callback.__name__ if self.callback else ""
class PriorityEventEmitter:
    """带优先级的通知系统"""
    def __init__(self):
        self._observers: Dict[str, List[PriorityObserver]] = {}
    def on(self, event: str, callback: Callable, priority: int = 0):
        """订阅事件(支持优先级)"""
        if event not in self._observers:
            self._observers[event] = []
        observer = PriorityObserver(priority, callback)
        heapq.heappush(self._observers[event], observer)
        print(f"[订阅] 事件 '{event}' (优先级:{priority})")
    def emit(self, event: str, *args, **kwargs):
        """按优先级触发事件"""
        if event not in self._observers:
            return
        # 按优先级排序执行
        observers = sorted(self._observers[event], key=lambda x: x.priority, reverse=True)
        for observer in observers:
            print(f"[执行] 优先级:{observer.priority}")
            observer.callback(*args, **kwargs)
# 使用示例 - 优先级通知
def high_priority_handler(msg):
    print(f"[高优先级] {msg}")
def medium_priority_handler(msg):
    print(f"[中优先级] {msg}")
def low_priority_handler(msg):
    print(f"[低优先级] {msg}")
# 测试
prioritized = PriorityEventEmitter()
prioritized.on("notification", low_priority_handler, -10)
prioritized.on("notification", high_priority_handler, 100)
prioritized.on("notification", medium_priority_handler, 0)
print("\n=== 优先级通知测试 ===")
prioritized.emit("notification", "系统更新完成")

异步事件通知系统

import asyncio
from typing import Any, Callable, Dict, List
class AsyncEventEmitter:
    """异步事件发射器"""
    def __init__(self):
        self._observers: Dict[str, List[Callable]] = {}
        self._loop = asyncio.get_event_loop()
    def on(self, event: str, callback: Callable):
        """订阅异步事件"""
        if event not in self._observers:
            self._observers[event] = []
        self._observers[event].append(callback)
    async def emit_async(self, event: str, *args, **kwargs):
        """异步触发事件"""
        if event not in self._observers:
            return
        tasks = []
        for callback in self._observers[event]:
            if asyncio.iscoroutinefunction(callback):
                tasks.append(callback(*args, **kwargs))
            else:
                callback(*args, **kwargs)
        if tasks:
            await asyncio.gather(*tasks, return_exceptions=True)
    def emit(self, event: str, *args, **kwargs):
        """同步触发异步事件"""
        asyncio.create_task(self.emit_async(event, *args, **kwargs))
# 异步观察者示例
async def async_logger(data):
    await asyncio.sleep(1)  # 模拟异步操作
    print(f"[异步日志] {data}")
async def async_email_sender(data):
    await asyncio.sleep(2)  # 模拟异步操作
    print(f"[异步邮件] 已发送: {data}")
async def main():
    emitter = AsyncEventEmitter()
    emitter.on("async_event", async_logger)
    emitter.on("async_event", async_email_sender)
    print("=== 异步通知测试 ===")
    await emitter.emit_async("async_event", {"message": "异步测试"})
    print("所有异步通知完成")
# 运行异步测试
if __name__ == "__main__":
    asyncio.run(main())

通用解决方案 - 装饰器模式

from functools import wraps
from typing import Any, Callable, Dict, List
class EventSystem:
    """完整的事件系统 - 支持装饰器"""
    def __init__(self):
        self.events: Dict[str, List[Callable]] = {}
    def observe(self, event: str):
        """装饰器:将函数注册为观察者"""
        def decorator(func: Callable):
            if event not in self.events:
                self.events[event] = []
            self.events[event].append(func)
            @wraps(func)
            def wrapper(*args, **kwargs):
                return func(*args, **kwargs)
            return wrapper
        return decorator
    def notify(self, event: str, *args, **kwargs):
        """通知所有观察者"""
        if event not in self.events:
            print(f"[通知] 事件 '{event}' 无观察者")
            return
        print(f"[通知] 事件 '{event}' 触发,通知 {len(self.events[event])} 个观察者")
        results = []
        for observer in self.events[event]:
            try:
                result = observer(*args, **kwargs)
                results.append(result)
            except Exception as e:
                print(f"[错误] 观察者 {observer.__name__} 执行失败: {e}")
        return results
# 使用装饰器
event_system = EventSystem()
@event_system.observe("user_login")
def send_welcome_email(username):
    msg = f"欢迎邮件已发送给 {username}"
    print(f"[观察者1] {msg}")
    return msg
@event_system.observe("user_login")
def create_user_session(username):
    msg = f"用户会话已创建: {username}"
    print(f"[观察者2] {msg}")
    return msg
@event_system.observe("user_login")
def update_last_login(username):
    msg = f"最后登录时间已更新: {username}"
    print(f"[观察者3] {msg}")
    return msg
# 测试
print("=== 装饰器模式测试 ===")
results = event_system.notify("user_login", "Alice")
print(f"所有观察者返回结果: {results}")

实际应用示例

# 完整的事件驱动系统
class DataEventSystem(EventSystem):
    """数据处理事件系统"""
    def __init__(self):
        super().__init__()
        self.data_store = {}
    def on_data_received(self, data):
        """处理接收到的数据"""
        print(f"\n[数据处理] 收到新数据: {data}")
        self.data_store[data['id']] = data
        self.notify("data_received", data)
    def on_data_error(self, error):
        """处理错误"""
        print(f"\n[错误处理] 错误: {error}")
        self.notify("data_error", error)
# 定义观察者
data_system = DataEventSystem()
@data_system.observe("data_received")
def process_data(data):
    print(f"[数据处理器] 处理数据: {data['name']}")
    # 模拟数据处理
    processed = f"处理结果: {data['value'] * 2}"
    print(f"  -> {processed}")
@data_system.observe("data_received")  
def log_data(data):
    print(f"[日志记录器] 记录数据: ID={data['id']}, Name={data['name']}")
@data_system.observe("data_error")
def handle_error(error):
    print(f"[错误处理器] 处理错误: {error}")
    # 发送告警等
# 模拟数据流
test_data = [
    {"id": 1, "name": "传感器A", "value": 42},
    {"id": 2, "name": "传感器B", "value": 78},
]
print("=== 实际应用示例 ===")
for data in test_data:
    data_system.on_data_received(data)
# 模拟错误
print("\n=== 触发错误事件 ===")
data_system.on_data_error("连接超时")

使用建议

  1. 选择合适模式

    • 简单通知:使用基础发布订阅
    • 需要优先级:使用优先级通知
    • 异步操作:使用异步事件系统
  2. 性能考虑

    • 观察者数量不多时,直接列表存储
    • 大量观察者时,考虑使用线程池或异步
  3. 错误处理

    • 单个观察者失败不应影响其他观察者
    • 添加超时和重试机制

这些实现可以根据你的具体需求进行调整和扩展,需要更具体的功能吗?

抱歉,评论功能暂时关闭!