Python异步可扩展数据增强场景适配设计模式问询
Hey there! Let's work through this problem together—your scenario has specific needs (async data enrichment, shared state between Event Stations (ES), easy extensibility) that the initial inheritance and facade approaches weren’t fully addressing. Here’s a tailored solution that checks all your boxes: a combination of the Composite Pattern with a Shared Context object, plus async-aware state management.
This approach keeps each ES focused on its single responsibility, lets you easily add new ES components, and safely manages shared data across async tasks.
1. Create a Shared Context Class
First, we’ll build a dedicated class to handle shared data between ES instances. It includes async locks to ensure thread-safe access in your async environment:
import asyncio from typing import Dict, Any class SharedContext: def __init__(self): self._shared_data: Dict[str, Any] = {} self._lock = asyncio.Lock() # Async-safe lock for concurrent access async def set_value(self, key: str, value: Any): async with self._lock: self._shared_data[key] = value async def get_value(self, key: str) -> Any: async with self._lock: return self._shared_data.get(key) async def update_values(self, update_dict: Dict[str, Any]): async with self._lock: self._shared_data.update(update_dict)
2. Refactor ES Classes for Single Responsibility
Each ES will now only handle its specific data enrichment logic, and rely on the SharedContext to read/write shared data. No more inheritance coupling! For example, here’s how ES3A would look:
# event_stations/es3a.py class ES3A: def __init__(self, context: SharedContext): self.context = context self.cache_table_name = None self.ignite_host = None self.ignite_port = None async def enrich_data(self, input_data: str): # Fetch shared metadata from the context metadata = await self.context.get_value("es_metadata") # Run your specific data enrichment logic here enhanced_result = f"ES3A_enhanced_{input_data}_{metadata.get('es3a_config')}" # Optional: Write results back to the context for other ES to use await self.context.set_value("es3a_result", enhanced_result) return enhanced_result
Repeat this pattern for all your ES classes (ES3B, ES3C, etc.)—each should only care about its own enrichment task and interact with shared data via the context.
3. Build a Coordinator Class (Replace Inheritance/Facade)
This class will manage ES initialization, load configuration, handle async task execution, and act as the central hub for your server logic. It follows the Open/Closed Principle—adding new ES only requires updating a list, not core code.
import os import asyncio import psycopg2 import websockets from datetime import datetime from websockets.extensions import permessage_deflate from structure import Structure from event_stations.es3a import ES3A from event_stations.es3b import ES3B from event_stations.es3c import ES3C from event_stations.es3d import ES3D from event_stations.es1a import ES1A from event_stations.es1b import ES1B from event_stations.es2a import ES2A from event_stations.es2b import ES2B from shared_context import SharedContext class FR_Server(Structure): unique_id = "100" fr_config_table = 'detail_event.app_fer_config' def __init__(self): print("Receiver Called INIT") self.psql_db = self.connect_to_psql() self.context = SharedContext() self.event_stations = self._initialize_event_stations() self._load_config_to_context() def connect_to_psql(self): return psycopg2.connect("dbname=trimble user=postgres password=admin") def _initialize_event_stations(self) -> list: # Initialize all ES instances with the shared context return [ ES1A(self.context), ES2A(self.context), ES2B(self.context), ES3A(self.context), ES3B(self.context), ES3C(self.context), ES3D(self.context) ] def _load_config_to_context(self): # Fetch config from DB, store in context, and set ES properties cur = self.psql_db.cursor() cur.execute(f"Select * from {self.fr_config_table}") config_results = cur.fetchall() es_metadata = {} for es_id, cache_table, ignite_port, ignite_host in config_results: es_metadata[es_id] = { "cache_table_name": cache_table, "ignite_port": ignite_port, "ignite_host": ignite_host } # Assign config to the matching ES instance for es in self.event_stations: if es.__class__.__name__ == es_id: es.cache_table_name = cache_table es.ignite_port = ignite_port es.ignite_host = ignite_host # Push metadata to the shared context for all ES to access asyncio.run(self.context.set_value("es_metadata", es_metadata)) def get_port(self): return os.getenv('WS_PORT', '10011') def get_host(self): return os.getenv('WS_HOST', 'localhost') async def start(self): return await websockets.serve( self.handler, self.get_host(), self.get_port(), ping_interval=None, max_size=None, max_queue=None, close_timeout=None, extensions=[ permessage_deflate.ServerPerMessageDeflateFactory( server_max_window_bits=11, client_max_window_bits=11, compress_settings={'memLevel': 4}, ), ] ) def generate_event_id(self, index): now = datetime.now() return "".join([ f"{now.day:02d}{now.month:02d}{now.year}{now.hour}{now.minute}{now.second}{now.microsecond}", self.unique_id, index ]) async def handler(self, websocket, path): async with websockets.connect( 'ws://localhost:10015', ping_interval=None, max_size=None, max_queue=None, close_timeout=None, extensions=[ permessage_deflate.ClientPerMessageDeflateFactory( server_max_window_bits=11, client_max_window_bits=11, compress_settings={'memLevel': 4}, ), ] ) as websocket_rb: async for row in websocket: lst_row = row.decode().split(",") uid = self.generate_event_id(lst_row[0]) lst_row = [uid] + lst_row # Create async tasks for all ES enrichment jobs enrichment_tasks = [es.enrich_data(lst_row[1]) for es in self.event_stations] # Run tasks in parallel and collect results results = await asyncio.gather(*enrichment_tasks) await websocket_rb.send(str(lst_row + results).encode()) def start_receiver(self): asyncio.get_event_loop().run_until_complete(self.start()) asyncio.get_event_loop().run_forever()
4. Key Benefits of This Approach
- Single Responsibility: Each ES handles only its data enrichment logic; the coordinator manages setup and scheduling; the context handles shared state.
- Easy Extensibility: To add a new ES, just:
- Implement the new ES class with a
enrich_dataasync method and context dependency - Add it to the
_initialize_event_stationslist - Add its config to the database table (if needed)
- Implement the new ES class with a
- Safe Async Shared State: The
SharedContextuses async locks to prevent race conditions when multiple ES access shared data. - Async Efficiency: Retains
asyncio.gatherfor parallel execution of enrichment tasks, keeping your pipeline fast.
内容的提问来源于stack exchange,提问作者Chetan Zope

