信号扫描模块设计与实现:多指标并行检测

发布时间:2026/8/25 21:06:00
信号扫描模块设计与实现:多指标并行检测 信号扫描模块设计与实现多指标并行检测行情数据是静态的但机会是动态的。手动盯盘在3个标的以内还能应付一旦超过10个人的注意力就成了瓶颈。信号扫描模块就是干这个的让代码替我们盯着所有标的一旦技术指标满足预设条件立刻输出信号。这篇文章分享一个实战可用的信号扫描模块核心设计思路多标的并行扫描ThreadPoolExecutor实现指标计算与信号判定分离规则可配置新增策略不用改扫描主逻辑1. 整体架构与扫描流程模块分三层数据层 - 指标层 - 信号层数据层负责拉取K线数据统一格式pd.DataFrame列名统一指标层计算均线、RSI、MACD输出指标DataFrame信号层根据规则判定是否触发信号输出信号列表扫描流程初始化标的列表 - 加载历史数据 - 计算技术指标 - 应用信号规则 - 输出结果关键设计决策并行扫描。每个标的的指标计算和信号判定是独立的天然适合并行。用concurrent.futures.ThreadPoolExecutor做并发I/O密集数据拉取和CPU密集指标计算都能受益。2. 数据层统一数据接口所有标的的数据格式必须一致否则指标计算代码没法复用。定义一个简单的数据加载函数import pandas as pd import numpy as np from typing import Dict, List, Optional import logging from concurrent.futures import ThreadPoolExecutor, as_completed logging.basicConfig(levellogging.INFO, format%(asctime)s - %(levelname)s - %(message)s) logger logging.getLogger(__name__) def load_kline_data(symbol: str, period: str 1d, limit: int 200) - Optional[pd.DataFrame]: 加载K线数据统一列名和格式。 这里用模拟数据演示实际场景替换为交易所API、数据库或本地文件。 # 模拟数据生成随机价格序列 np.random.seed(hash(symbol) % 2**32) dates pd.date_range(endpd.Timestamp.today(), periodslimit, freqperiod) base_price 100 np.random.rand() * 50 prices base_price * (1 np.random.randn(limit).cumsum() * 0.02) df pd.DataFrame({ date: dates, open: prices * (1 np.random.randn(limit) * 0.005), high: prices * (1 np.abs(np.random.randn(limit)) * 0.01), low: prices * (1 - np.abs(np.random.randn(limit)) * 0.01), close: prices, volume: np.random.randint(10000, 100000, sizelimit) }) df.set_index(date, inplaceTrue) return df实际项目中这里通常对接ccxt、akshare或本地数据库。关键是返回的DataFrame必须包含open/high/low/close/volume五列索引是时间。3. 指标层计算技术指标指标计算用pandas向量化操作性能好且代码简洁。核心指标均线MA、RSI、MACD。def calculate_indicators(df: pd.DataFrame) - pd.DataFrame: 计算技术指标返回包含原始数据和指标列的DataFrame。 if df is None or df.empty: return df result df.copy() # 均线系统MA5, MA10, MA20, MA60 for window in [5, 10, 20, 60]: result[fma_{window}] result[close].rolling(windowwindow).mean() # RSI (14日) delta result[close].diff() gain delta.where(delta 0, 0.0) loss -delta.where(delta 0, 0.0) avg_gain gain.rolling(window14).mean() avg_loss loss.rolling(window14).mean() rs avg_gain / avg_loss result[rsi_14] 100 - (100 / (1 rs)) # MACD (12, 26, 9) ema_12 result[close].ewm(span12, adjustFalse).mean() ema_26 result[close].ewm(span26, adjustFalse).mean() result[macd_dif] ema_12 - ema_26 result[macd_dea] result[macd_dif].ewm(span9, adjustFalse).mean() result[macd_hist] result[macd_dif] - result[macd_dea] return result注意几个细节rolling(window).mean()计算移动平均前N-1个值是NaNRSI 计算里delta.where(delta 0, 0.0)把负变动置0避免负数干扰MACD 用ewm指数加权比SMA更贴近实际交易软件的计算方式4. 信号层规则判定信号判定是整个模块的核心。设计思路规则函数接收指标DataFrame返回布尔值True触发信号。这样每种策略就是一个独立函数互不干扰。def check_ma_cross(df: pd.DataFrame) - bool: 金叉信号MA5上穿MA20 if len(df) 20: return False latest df.iloc[-1] prev df.iloc[-2] # 前一根K线 MA5 MA20当前K线 MA5 MA20 return (prev[ma_5] prev[ma_20] and latest[ma_5] latest[ma_20]) def check_rsi_oversold(df: pd.DataFrame) - bool: 超卖反弹信号RSI从低于30回升到30以上 if len(df) 15: return False latest df.iloc[-1] prev df.iloc[-2] return (prev[rsi_14] 30 and latest[rsi_14] 30) def check_macd_golden_cross(df: pd.DataFrame) - bool: MACD金叉DIF上穿DEA if len(df) 26: return False latest df.iloc[-1] prev df.iloc[-2] return (prev[macd_dif] prev[macd_dea] and latest[macd_dif] latest[macd_dea])每个函数只做一件事判断当前时刻是否满足信号条件。注意边界处理数据长度不够时直接返回False避免 NaN 比较出错。5. 并行扫描主逻辑扫描模块的核心对每个标的执行加载数据 - 计算指标 - 跑所有信号规则。class SignalScanner: 信号扫描器 def __init__(self, symbols: List[str], signal_rules: Dict[str, callable], max_workers: int 4): self.symbols symbols self.signal_rules signal_rules # {规则名: 函数} self.max_workers max_workers def scan_symbol(self, symbol: str) - Dict: 扫描单个标的返回信号结果 try: # 1. 加载数据 df load_kline_data(symbol) if df is None or len(df) 60: return {symbol: symbol, signals: [], error: 数据不足} # 2. 计算指标 df_indicators calculate_indicators(df) # 3. 运行所有信号规则 triggered [] for rule_name, rule_func in self.signal_rules.items(): try: if rule_func(df_indicators): triggered.append(rule_name) except Exception as e: logger.warning(f[{symbol}] 规则 {rule_name} 执行异常: {e}) return {symbol: symbol, signals: triggered, error: None} except Exception as e: logger.error(f[{symbol}] 扫描异常: {e}) return {symbol: symbol, signals: [], error: str(e)} def scan_all(self) - List[Dict]: 并行扫描所有标的 results [] with ThreadPoolExecutor(max_workersself.max_workers) as executor: # 提交所有任务 future_map {executor.submit(self.scan_symbol, sym): sym for sym in self.symbols} # 收集结果 for future in as_completed(future_map): symbol future_map[future] try: result future.result() results.append(result) # 打印日志 if result[signals]: logger.info(f[{symbol}] 触发信号: {result[signals]}) except Exception as e: logger.error(f[{symbol}] 任务执行失败: {e}) return results关键点ThreadPoolExecutor控制并发数避免请求过多被交易所限流as_completed按完成顺序处理结果不用等所有任务结束单个标的异常不影响其他标的的扫描信号规则用字典传入新增策略只需加一个函数和一个键值对6. 使用示例if __name__ __main__: # 定义要扫描的标的 symbols [BTC_USDT, ETH_USDT, BNB_USDT, SOL_USDT, XRP_USDT] # 注册信号规则 rules { MA金叉: check_ma_cross, RSI超卖反弹: check_rsi_oversold, MACD金叉: check_macd_golden_cross, } # 创建扫描器并执行 scanner SignalScanner(symbols, rules, max_workers3) results scanner.scan_all() # 输出结果 print(\n 扫描结果 ) for res in results: status ✅ if res[signals] else ❌ print(f{status} {res[symbol]}: {res[signals] if res[signals] else 无信号})输出效果2024-11-20 10:30:01 - INFO - [BTC_USDT] 触发信号: [MA金叉, MACD金叉] 2024-11-20 10:30:02 - INFO - [ETH_USDT] 触发信号: [RSI超卖反弹] 2024-11-20 10:30:02 - INFO - [SOL_USDT] 触发信号: [MA金叉] 扫描结果 ✅ BTC_USDT: [MA金叉, MACD金叉] ❌ ETH_USDT: [RSI超卖反弹] ✅ BNB_USDT: [] ...7. 性能与扩展性优化性能瓶颈指标计算是纯CPU操作数据拉取是I/O操作。ThreadPoolExecutor对I/O密集型任务提升明显但对纯CPU计算受GIL限制。如果扫描几百个标的且数据量大可以考虑ProcessPoolExecutor做CPU密集型指标计算数据层加缓存functools.lru_cache或 Redis避免重复拉取增量更新只计算最新K线而不是每次全量重算扩展方向信号规则组合目前每个规则独立判定可以增加组合逻辑如MA金叉 AND RSI50信号优先级多个信号同时触发时按权重排序输出实时扫描配合定时器或WebSocket做到分钟级扫描结果持久化把信号写入数据库或消息队列供下游策略模块消费一个比较实用的增量更新示例def incremental_indicators(df: pd.DataFrame, last_indicators: Dict) - pd.DataFrame: 增量计算指标传入上次的指标结果只更新最新数据。 大幅减少计算量适合高频扫描场景。 # 简化版只保留最近200根K线 df df.tail(200) return calculate_indicators(df)8. 踩坑记录写这个模块时遇到的几个实际问题NaN 比较陷阱rolling计算前N-1个值是NaN直接比较会报错。处理方式数据长度不足时提前返回False。时区问题不同交易所返回的K线时间戳时区不一致统一转UTC再处理。并发数控制线程数不是越多越好交易所API有频率限制。建议并发数 标的数开根号或者根据API限制动态调整。异常隔离一个标的的数据异常如停牌、数据缺失不应该影响其他标的的扫描。用try-except包裹每个标的的处理逻辑。9. 结语信号扫描模块的核心价值在于把人工盯盘的经验转化为可重复执行的代码逻辑。通过并行扫描把注意力从盯盘解放出来留给真正的策略优化和风控。这个模块的设计是通用的指标公式和信号规则都可以按需替换。后续可以在此基础上扩展信号历史统计、回测验证、多时间框架扫描等。更多内容请关注本站。