解密OpenHands插件系统:手把手教你开发自定义微代理
OpenHands插件系统深度开发指南:构建电商价格监控微代理实战
1. 理解OpenHands插件系统的核心架构
OpenHands的插件系统是其最具扩展性的设计之一,它允许开发者在不修改核心代码的情况下,为平台添加新功能。这套系统基于Python的模块化架构,通过精心设计的接口与核心组件交互。
核心机制解析:
- 动态加载:插件在运行时通过
importlib动态加载,无需重启主程序 - 接口契约:所有插件必须实现
BasePlugin抽象类定义的接口方法 - 依赖隔离:每个插件运行在独立的虚拟环境中,避免依赖冲突
- 热插拔:支持插件的安装、卸载和更新而不影响系统稳定性
插件系统与OpenHands其他模块的交互关系如下表所示:
| 系统模块 | 插件接入点 | 典型应用场景 |
|---|---|---|
| 事件系统 | 事件订阅/发布 | 实时响应价格变动事件 |
| 工具系统 | 工具注册 | 添加自定义数据采集器 |
| 内存系统 | 上下文注入 | 持久化监控配置 |
| 运行时 | 环境钩子 | 浏览器自动化控制 |
提示:开发前建议先熟悉
plugin_interface.py中的类型定义,特别是PluginMetadata和PluginConfig这两个数据类
2. 搭建开发环境与工具链配置
2.1 环境准备
推荐使用Python 3.10+版本,并创建独立的虚拟环境:
python -m venv openhands-dev source openhands-dev/bin/activate # Linux/Mac # 或 openhands-dev\Scripts\activate # Windows安装核心依赖包:
pip install openhands-core>=1.2.0 pip install playwright aiohttp beautifulsoup42.2 项目结构规划
电商价格监控插件建议采用以下目录结构:
ecommerce_monitor/ ├── __init__.py ├── plugin.py # 主插件类 ├── config_schema.py # 配置模型 ├── utils/ # 工具函数 │ ├── scraper.py │ └── notifier.py └── tests/ # 单元测试 └── test_plugin.py关键文件plugin.py的基本框架:
from dataclasses import dataclass from typing import Optional from openhands.plugins import BasePlugin @dataclass class PriceMonitorConfig: target_url: str check_interval: int = 3600 price_threshold: Optional[float] = None class EcommerceMonitorPlugin(BasePlugin): def __init__(self, config: PriceMonitorConfig): self.config = config self._is_running = False async def start(self): """启动监控任务""" self._is_running = True # 实现具体逻辑 async def stop(self): """停止监控""" self._is_running = False3. 实现电商价格监控的核心逻辑
3.1 网页内容抓取策略
现代电商网站通常采用动态加载技术,需要结合多种采集方式:
API直接调用(最优选择)
- 通过浏览器开发者工具分析XHR请求
- 模拟合法请求头避免被封禁
浏览器自动化(适用于复杂SPA)
from playwright.async_api import async_playwright async def fetch_dynamic_price(url): async with async_playwright() as p: browser = await p.chromium.launch() page = await browser.new_page() await page.goto(url) # 等待价格元素加载 await page.wait_for_selector('.price-section') # 执行页面内JS获取数据 price = await page.evaluate('''() => { return parseFloat( document.querySelector('.final-price').innerText .replace('$','').trim() ) }''') await browser.close() return price静态HTML解析(轻量级方案)
from bs4 import BeautifulSoup import aiohttp async def fetch_static_price(url): async with aiohttp.ClientSession() as session: async with session.get(url) as resp: html = await resp.text() soup = BeautifulSoup(html, 'html.parser') price_str = soup.find(class_='price').get_text() return float(price_str[1:].replace(',',''))
3.2 价格变化检测算法
简单的阈值比较可能产生大量误报,推荐采用滑动窗口均值算法:
from collections import deque import statistics class PriceChangeDetector: def __init__(self, window_size=5, sensitivity=0.1): self.window = deque(maxlen=window_size) self.sensitivity = sensitivity def check_abnormal(self, current_price): if not self.window: self.window.append(current_price) return False avg = statistics.mean(self.window) change_rate = abs(current_price - avg) / avg self.window.append(current_price) return change_rate > self.sensitivity3.3 事件订阅与通知系统
将价格变动作为自定义事件发布到OpenHands事件总线:
from openhands.events import Event, EventStream class PriceDropEvent(Event): def __init__(self, product_url, old_price, new_price): self.product_url = product_url self.old_price = old_price self.new_price = new_price async def monitor_loop(plugin): detector = PriceChangeDetector() while plugin._is_running: current_price = await fetch_price(plugin.config.target_url) if detector.check_abnormal(current_price): event = PriceDropEvent( plugin.config.target_url, detector.window[-2], current_price ) await plugin.event_stream.publish(event) await asyncio.sleep(plugin.config.check_interval)4. 插件集成与高级功能扩展
4.1 注册插件工具
将价格监控功能暴露为系统工具,供其他代理调用:
from openhands.mcp import Tool, ToolParameter class PriceCheckTool(Tool): name = "ecommerce_price_check" description = "获取指定商品当前价格" parameters = [ ToolParameter( name="product_url", type="string", description="商品页面URL", required=True ) ] async def execute(self, product_url: str): return await fetch_price(product_url) # 在插件启动时注册 def on_enable(self): self.register_tool(PriceCheckTool())4.2 持久化配置管理
利用OpenHands的存储系统保存监控记录:
async def save_price_history(self, price_data): storage = self.get_storage() history = await storage.get('price_history') or [] history.append({ 'timestamp': datetime.now().isoformat(), 'price': price_data }) await storage.set('price_history', history[-1000:]) # 保留最近1000条4.3 可视化监控面板
通过Web扩展点添加管理界面:
from fastapi import APIRouter from fastapi.responses import HTMLResponse router = APIRouter() @router.get("/monitor-dashboard", response_class=HTMLResponse) async def get_dashboard(): return """ <html> <head><title>价格监控面板</title></head> <body> <div id="price-chart"></div> <script> // 使用Chart.js渲染价格曲线 </script> </body> </html> """ def get_web_routers(self): return [router]5. 调试与性能优化技巧
5.1 常见问题排查指南
| 问题现象 | 可能原因 | 解决方案 |
|---|---|---|
| 获取价格超时 | 网站反爬机制触发 | 1. 增加请求间隔 2. 轮换User-Agent 3. 使用代理IP池 |
| 价格解析失败 | 页面结构变更 | 1. 更新CSS选择器 2. 添加多重fallback解析逻辑 |
| 事件未被处理 | 订阅者未正确注册 | 1. 检查事件类型匹配 2. 验证订阅回调函数签名 |
5.2 性能优化策略
浏览器实例复用优化:
class BrowserManager: _instance = None @classmethod async def get_browser(cls): if cls._instance is None: cls._p = await async_playwright().start() cls._instance = await cls._p.chromium.launch( headless=True, args=['--disable-blink-features=AutomationControlled'] ) return cls._instance # 使用方式 browser = await BrowserManager.get_browser() page = await browser.new_page()智能请求调度算法:
class RequestScheduler: def __init__(self, max_concurrent=3): self.semaphore = asyncio.Semaphore(max_concurrent) self.domain_delay = defaultdict(float) async def fetch(self, url): domain = urlparse(url).netloc now = time.time() # 遵守域名的请求间隔限制 if now < self.domain_delay[domain]: await asyncio.sleep(self.domain_delay[domain] - now) async with self.semaphore: try: result = await do_fetch(url) self.domain_delay[domain] = time.time() + 2 # 2秒间隔 return result except Exception: self.domain_delay[domain] = time.time() + 60 # 出错时延长间隔 raise6. 生产环境部署方案
6.1 容器化打包
创建优化的Docker镜像:
FROM python:3.10-slim WORKDIR /app COPY requirements.txt . RUN pip install --no-cache-dir -r requirements.txt \ && playwright install chromium \ && playwright install-deps COPY . . CMD ["openhands", "run-plugin", "ecommerce_monitor"]构建命令:
docker build -t openhands-monitor . docker run -d \ -e CONFIG_PATH=/app/config.yaml \ -v ./config:/app/config \ openhands-monitor6.2 配置管理最佳实践
推荐采用分层配置:
# config.yaml default: check_interval: 3600 notification_methods: - email - webhook products: - url: "https://example.com/product1" price_threshold: 99.99 selectors: price: ".sale-price" - url: "https://example.com/product2" check_interval: 1800 # 覆盖全局设置加载配置的代码实现:
import yaml from pydantic import BaseModel class ProductConfig(BaseModel): url: str price_threshold: Optional[float] selectors: dict def load_config(path): with open(path) as f: data = yaml.safe_load(f) base_config = data['default'] products = [ProductConfig(**p) for p in data['products']] return base_config, products7. 插件生态系统集成
7.1 与现有微代理协同工作
通过事件总线实现插件间的松耦合协作:
async def handle_inventory_event(event: InventoryUpdateEvent): if event.stock_status == 'in_stock': # 当库存恢复时触发价格检查 price = await PriceCheckTool().execute(event.product_url) if price < self.config.price_threshold: await notify_price_drop(event.product_url, price) def on_enable(self): self.event_stream.subscribe( 'inventory_update', handle_inventory_event )7.2 构建插件组合工作流
利用OpenHands的管道特性串联多个插件:
from openhands.workflow import Pipeline pipeline = Pipeline( name='price_monitor_workflow', steps=[ ('product_finder', {'category': 'electronics'}), ('price_monitor', {'check_interval': 1800}), ('alert_sender', {'channels': ['slack']}) ] ) await pipeline.run()