本文目录导读:

我来帮你实现一个事件多观察者通知系统,这里提供几种不同场景的实现方案:
基础发布订阅模式
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("连接超时")
使用建议
-
选择合适模式:
- 简单通知:使用基础发布订阅
- 需要优先级:使用优先级通知
- 异步操作:使用异步事件系统
-
性能考虑:
- 观察者数量不多时,直接列表存储
- 大量观察者时,考虑使用线程池或异步
-
错误处理:
- 单个观察者失败不应影响其他观察者
- 添加超时和重试机制
这些实现可以根据你的具体需求进行调整和扩展,需要更具体的功能吗?