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

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.

Solution: Composite Pattern + Shared Context

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:
    1. Implement the new ES class with a enrich_data async method and context dependency
    2. Add it to the _initialize_event_stations list
    3. Add its config to the database table (if needed)
  • Safe Async Shared State: The SharedContext uses async locks to prevent race conditions when multiple ES access shared data.
  • Async Efficiency: Retains asyncio.gather for parallel execution of enrichment tasks, keeping your pipeline fast.

内容的提问来源于stack exchange,提问作者Chetan Zope

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 21:37:51