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

Python动态增删模块:无需重启新增加密货币交易所数据抓取

Dynamic Exchange Integration Without Restarting Your Crypto Data Scraper

Great question! This is a perfect scenario for building an extensible, plugin-based architecture that lets you add new exchange support on the fly—all while keeping your existing Binance/Poloniex scrapers running uninterrupted. Let’s break down the approach with concrete steps and code examples (I’ll use Python since it’s widely used for scraping projects like this).

1. Start with a Unified Exchange Interface

First, define an abstract base class (ABC) that standardizes how all exchange fetchers behave. This ensures every new exchange you add follows the same rules, making dynamic loading trivial.

from abc import ABC, abstractmethod
from typing import Dict, Any
import threading
import time

class ExchangeFetcher(ABC):
    def __init__(self, config: Dict[str, Any]):
        self.config = config
        self.running = True
        self.thread = threading.Thread(target=self._run_loop, daemon=True)

    @abstractmethod
    def fetch_market_data(self) -> Dict[str, Any]:
        """Fetch market data (prices, volumes, etc.) from the exchange"""
        pass

    @abstractmethod
    def get_exchange_name(self) -> str:
        """Return the unique name of the exchange (e.g., 'binance')"""
        pass

    def _run_loop(self):
        """Internal loop to fetch data at the configured interval"""
        interval = self.config.get("fetch_interval", 10)
        while self.running:
            try:
                data = self.fetch_market_data()
                # Process/store the data (e.g., send to a database, Kafka topic)
                print(f"Fetched data from {self.get_exchange_name()}: {data}")
            except Exception as e:
                print(f"Error fetching from {self.get_exchange_name()}: {str(e)}")
            time.sleep(interval)

    def start(self):
        """Start the fetching thread"""
        self.thread.start()

    def stop(self):
        """Gracefully stop the fetching thread"""
        self.running = False
        self.thread.join()

2. Adapt Existing Exchanges to the Interface

Update your Binance and Poloniex scrapers to inherit from this interface. Here’s a simplified Binance example:

import requests

class BinanceFetcher(ExchangeFetcher):
    def fetch_market_data(self) -> Dict[str, Any]:
        response = requests.get("https://api.binance.com/api/v3/ticker/24hr")
        response.raise_for_status()
        return {"tickers": response.json()[:5]}  # Return first 5 tickers for example

    def get_exchange_name(self) -> str:
        return "binance"

Do the same for Poloniex—you’ll just swap out the API endpoint and parsing logic to match their API structure.

3. Build a Dynamic Loader for New Exchanges

Now, create a manager class that handles registering, starting, and dynamically loading new fetchers. You have two practical options for dynamic loading:

Option A: Monitor a "Plugins" Directory

Use a library like watchdog to watch a dedicated directory for new Python files (e.g., kraken_fetcher.py). When a new file is added, the manager automatically loads it and starts the fetcher.

import importlib.util
import os
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler

class ExchangeManager:
    def __init__(self, plugin_dir: str = "./exchange_plugins"):
        self.plugin_dir = plugin_dir
        self.fetchers = {}
        # Create plugin dir if it doesn't exist
        os.makedirs(plugin_dir, exist_ok=True)
        # Load initial plugins (Binance/Poloniex)
        self._load_all_plugins()
        # Start watching for new plugins
        self._start_plugin_watcher()

    def _load_plugin(self, file_path: str):
        """Load a single exchange plugin from a file"""
        module_name = os.path.basename(file_path).replace(".py", "")
        spec = importlib.util.spec_from_file_location(module_name, file_path)
        module = importlib.util.module_from_spec(spec)
        spec.loader.exec_module(module)
        
        # Find the ExchangeFetcher subclass in the module
        for attr in dir(module):
            cls = getattr(module, attr)
            if isinstance(cls, type) and issubclass(cls, ExchangeFetcher) and cls != ExchangeFetcher:
                # Pull config from a dynamic store (e.g., Redis, env vars)
                config = {"fetch_interval": 10}
                fetcher = cls(config)
                self.register_fetcher(fetcher)
                break

    def _load_all_plugins(self):
        """Load all existing plugins in the plugin directory"""
        for file in os.listdir(self.plugin_dir):
            if file.endswith(".py") and not file.startswith("_"):
                self._load_plugin(os.path.join(self.plugin_dir, file))

    def _start_plugin_watcher(self):
        """Watch the plugin directory for new files"""
        class PluginHandler(FileSystemEventHandler):
            def on_created(self, event):
                if not event.is_directory and event.src_path.endswith(".py"):
                    self._load_plugin(event.src_path)

        handler = PluginHandler()
        handler._load_plugin = self._load_plugin  # Bind the manager's method
        observer = Observer()
        observer.schedule(handler, self.plugin_dir, recursive=False)
        observer.start()

    def register_fetcher(self, fetcher: ExchangeFetcher):
        """Register and start a new fetcher"""
        exchange_name = fetcher.get_exchange_name()
        if exchange_name not in self.fetchers:
            self.fetchers[exchange_name] = fetcher
            fetcher.start()
            print(f"Started fetching from {exchange_name}")

Option B: Add an API Endpoint to Trigger Loading

If you prefer to control loading via an API (e.g., for automated deployments), add a simple Flask/FastAPI route that accepts a plugin file path or module reference, then loads it dynamically.

from fastapi import FastAPI

app = FastAPI()
exchange_manager = ExchangeManager()

@app.post("/add-exchange")
def add_exchange(plugin_path: str):
    exchange_manager._load_plugin(plugin_path)
    return {"status": "success", "message": f"Loaded plugin from {plugin_path}"}

4. Implement the Kraken Fetcher

To add Kraken support, just create a kraken_fetcher.py file in your exchange_plugins directory:

import requests

class KrakenFetcher(ExchangeFetcher):
    def fetch_market_data(self) -> Dict[str, Any]:
        response = requests.get("https://api.kraken.com/0/public/Ticker", params={"pair": "XBTUSD,ETHUSD"})
        response.raise_for_status()
        return {"tickers": response.json()["result"]}

    def get_exchange_name(self) -> str:
        return "kraken"

As soon as you save this file, the ExchangeManager will detect it, load the class, and start fetching data from Kraken—no restart needed!

5. Key Production Considerations

  • Dynamic Configuration: Use a config store like Redis or etcd instead of hardcoding API keys or intervals. This lets you update settings without touching code.
  • Isolation: For critical systems, run each fetcher in a separate process (using multiprocessing) instead of a thread—so a bug in one exchange’s code doesn’t crash others.
  • Monitoring: Add metrics (e.g., Prometheus) to track fetch success rates, latency, and data volume per exchange.
  • Graceful Shutdown: Ensure the manager can stop all fetchers cleanly, and handle plugin reloads without losing in-flight data.

内容的提问来源于stack exchange,提问作者Martin Pichler

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 04:06:24