第 10 章 · 01 选股调度:一票否决与首选制 本节摘要:本节解剖 ABU 回测流水线的"选股调度中枢"—— + + ( 约 290 行选股部分)。 要对全市场几百只股票回测,但不是每只都值得回测——先要"选股"。 是对外接口:choicesymbols 为 None 时调 取全市场 → 只有 FORCELOCAL 模式才允许并行选股(HDF5 多进程易写坏,否则强制 nprocess=1) → 切分子序列 → 分发到多进程 → 合并结果 → 训练测试集切割。 是选股执行者: 用 dict 实例化各选股因子,实行一票否决制 + 首选制—— 的因子进 (层层过滤,批量高效),其余进 ( 一票否决,任一因子反对即 break 剔除)。
本节摘要:本节解剖 ABU 回测流水线的"选股调度中枢"——
AbuPickStockMaster+AbuPickStockWorker+ABuPickStockExecute(AlphaBu约 290 行选股部分)。run_loop_back要对全市场几百只股票回测,但不是每只都值得回测——先要"选股"。AbuPickStockMaster.do_pick_stock_with_process是对外接口:choice_symbols 为 None 时调all_symbol()取全市场 → 只有 FORCE_LOCAL 模式才允许并行选股(HDF5 多进程易写坏,否则强制 n_process=1) →split_k_market切分子序列 →Parallel+delayed分发到多进程 →itertools.chain合并结果 → 训练测试集切割。AbuPickStockWorker是选股执行者:init_stock_pickers用 dict 实例化各选股因子,实行一票否决制 + 首选制——first_choice:True的因子进first_stock_pickers(层层过滤,批量高效),其余进stock_pickers(_batch_fit一票否决,任一因子反对即 break 剔除)。本节会精读这两个机制的代码,以及为什么 ABU 选股并行受 HDF5 限制。
内容来源:原项目源码
abupy/AlphaBu/ABuPickBase.py(64 行)+ABuPickStockMaster.py(109 行)+ABuPickStockWorker.py(144 行)+ABuPickStockExecute.py(52 行)。
⚠️ 注意:本节重点理解"为什么选股并行受限"。ABU 的 KL 数据默认存 HDF5,HDF5 多进程写会损坏文件(并发写不加锁),所以选股阶段只有在
E_DATA_FETCH_FORCE_LOCAL(强制本地缓存)模式下才允许并行(此时只读不写),否则强制 n_process=1 单进程。这是 ABU 在"性能 vs 数据安全"间的保守取舍。另外"一票否决制"是 ABU 选股的核心哲学——任一因子反对就剔除,保证选出的股票被所有因子认可。
阅读完本节,你应当能够:
do_pick_stock_with_process:choice_symbols None 调 all_symbol → FORCE_LOCAL 才并行否则 n_process=1 → split_k_market 切分 → Parallel+delayed 分发 → itertools.chain 合并 → 训练测试集切割。do_pick_stock_work 的前后包装:AbuKLManager 构造 → AbuPickStockWorker 实例化 → fit → 返回 choice_symbols。init_stock_pickers:dict 实例化 → first_choice:True 进 first_stock_pickers(首选),其余进 stock_pickers(普通)。_first_batch_fit(首选,层层过滤 fit_first_choice)与 _batch_fit(普通,外层 symbol 内层因子,fit_pick 返回 False 就 break 剔除)的区别。ABuPickBase.py(64 行)定义两个抽象基类,都继承 AbuParamBase(ABU 的参数化基类,支持 dict 构造):
19 class AbuPickTimeWorkBase(six.with_metaclass(ABCMeta, AbuParamBase)): 22 @abstractmethod 23 def fit(self, *args, **kwargs): 31 @abstractmethod 32 def init_sell_factors(self, *args, **kwargs): 39 @abstractmethod 40 def init_buy_factors(self, *args, **kwargs): 46 class AbuPickStockWorkBase(six.with_metaclass(ABCMeta, AbuParamBase)): 49 @abstractmethod 50 def fit(self, *args, **kwargs): 58 @abstractmethod 59 def init_stock_pickers(self, *args, **kwargs):
两个基类的抽象方法对称:
fit(开始择时)+ init_sell_factors(初始化卖出因子)+ init_buy_factors(初始化买入因子)。择时是"针对一个 symbol 的时间序列,决定每天买还是卖"。fit(开始选股)+ init_stock_pickers(初始化选股因子)。选股是"针对多个 symbol,决定哪些值得纳入回测"。注意选股和择时的本质区别:选股是数据多对多(多个 symbol 对多个选股因子),择时是时间驱动(单个 symbol 的时间序列逐日推进)。这个区别决定了它们的 fit 实现完全不同(下节讲择时)。
do_pick_stock_with_process(ABuPickStockMaster.py:28-99)是对外接口,也是选股并行的总指挥:
28 @classmethod 29 def do_pick_stock_with_process(cls, capital, benchmark, stock_pickers, choice_symbols=None, 30 n_process_pick_stock=ABuEnv.g_cpu_cnt, callback=None): 42 input_choice_symbols = True 43 if choice_symbols is None or len(choice_symbols) == 0: 44 choice_symbols = all_symbol() # None 调 all_symbol 取全市场 45 input_choice_symbols = False 47 if n_process_pick_stock <= 0: 49 n_process_pick_stock = ABuEnv.g_cpu_cnt 50 if stock_pickers is not None: 54 if n_process_pick_stock > 1 and ABuEnv.g_data_fetch_mode != EMarketDataFetchMode.E_DATA_FETCH_FORCE_LOCAL: 57 logging.info('batch get only support E_DATA_FETCH_FORCE_LOCAL for Parallel!') 58 n_process_pick_stock = 1 # 非 FORCE_LOCAL 强制单进程 61 process_symbols = split_k_market(n_process_pick_stock, market_symbols=choice_symbols) 64 if n_process_pick_stock > 1: 65 n_process_pick_stock = len(process_symbols) # 切割余数,32->33 67 parallel = Parallel(n_jobs=n_process_pick_stock, verbose=0, pre_dispatch='2*n_jobs') 70 if callback is None: 71 callback = do_pick_stock_work 74 p_nev = AbuEnvProcess() 76 out_choice_symbols = parallel(delayed(callback)(choice_symbols, benchmark, capital, stock_pickers, env=p_nev) 79 for choice_symbols in process_symbols) 82 choice_symbols = list(itertools.chain.from_iterable(out_choice_symbols)) # 合并
choice_symbols is None 时调 all_symbol() 取全市场(美股几千只 / A 股几千只 / 比特币等)。input_choice_symbols = False 记录"用户没指定",后面训练测试集切割要用这个标志(只有全市场才切割,用户指定的不切)。
这是选股调度的关键限制:
54 if n_process_pick_stock > 1 and ABuEnv.g_data_fetch_mode != EMarketDataFetchMode.E_DATA_FETCH_FORCE_LOCAL: 57 logging.info('batch get only support E_DATA_FETCH_FORCE_LOCAL for Parallel!') 58 n_process_pick_stock = 1
源码注释(L55-56)给出两条原因:
所以默认(非 FORCE_LOCAL)强制 n_process=1 单进程。这是 ABU 在"性能 vs 数据安全"间的保守取舍——宁可慢,也不能写坏数据。
split_k_market(n_process, market_symbols) 把 symbols 按"市场+代码"排序后均匀切成 n_process 份。切割会有余数(比如 100 只切 3 份可能切出 4 份),所以 L64-65 把实际进程数调整为切割后的份数。
Parallel 是 joblib 的并行执行器(ABuParallel 封装),delayed(callback)(...) 把 do_pick_stock_work 包装成延迟任务。每个进程拿到一份 symbols,独立执行选股,返回筛选后的 symbols 列表。itertools.chain.from_iterable 把多个进程的结果扁平化合并。
选股完成后,根据 env 设置切割训练/测试集:
89 if not input_choice_symbols and ABuEnv.g_enable_last_split_test: 91 choice_symbols = ABuMarket.market_last_split_test() # 只用上次切割的测试集 92 elif not input_choice_symbols and ABuEnv.g_enable_last_split_train: 94 choice_symbols = ABuMarket.market_last_split_train() # 只用上次切割的训练集 95 elif ABuEnv.g_enable_train_test_split: 97 choice_symbols = ABuMarket.market_train_test_split(ABuEnv.g_split_tt_n_folds, choice_symbols) # 切割并返回训练集
三个开关有优先级:g_enable_last_split_test > g_enable_last_split_train > g_enable_train_test_split。注意只有 not input_choice_symbols(用户没指定)时才用前两个——用户指定的 symbols 不强制切割。训练测试集切割是 ABU 防止"用同一批股票既调参又回测"的过拟合保护(第 6 章度量详讲)。
ABuPickStockExecute.py:20-34 是 Worker 的前后包装:
20 @add_process_env_sig 21 def do_pick_stock_work(choice_symbols, benchmark, capital, stock_pickers): 30 kl_pd_manager = AbuKLManager(benchmark, capital) # 先建 KL 管理器 31 stock_pick = AbuPickStockWorker(capital, benchmark, kl_pd_manager, choice_symbols=choice_symbols, 32 stock_pickers=stock_pickers) # 实例化 Worker 33 stock_pick.fit() # 执行选股 34 return stock_pick.choice_symbols # 返回筛选后的 symbols
@add_process_env_sig 装饰器(ABuEnvProcess.py)处理进程间环境变量传递——多进程下子进程的环境需要从主进程拷贝,这个装饰器统一处理。三步:建 KL 管理器 → 实例化 Worker → fit → 返回 choice_symbols。
do_pick_stock_thread_work(L37-51)是已废弃的"进程+线程"混合模式(用 ThreadPoolExecutor),源码标记 Deprecated,因为 hdf5 多线程读取也有问题。ABU 最终选了纯多进程。
AbuPickStockWorker(ABuPickStockWorker.py:23-143)是选股的核心执行者。构造和因子初始化:
26 def __init__(self, capital, benchmark, kl_pd_manager, choice_symbols=None, stock_pickers=None): 38 self.stock_pickers = [] # 普通选股因子(一票否决) 39 self.first_stock_pickers = [] # 首选因子(层层过滤) 40 self.init_stock_pickers(stock_pickers) 48 def init_stock_pickers(self, stock_pickers): 54 if stock_pickers is not None: 55 for picker_class in stock_pickers: 63 picker_class_cp = copy.deepcopy(picker_class) 65 class_fac = picker_class_cp.pop('class') # pop 出 class, 剩下是参数 67 picker = class_fac(self.capital, self.benchmark, **picker_class_cp) # dict 实例化 73 if 'first_choice' in picker_class and picker_class['first_choice']: 74 self.first_stock_pickers.append(picker) # first_choice:True 进首选 76 else: 77 self.stock_pickers.append(picker) # 其余进普通
stock_pickers 是 dict 列表,每个 dict 形如 {'class': AbuPickStockNTop, 'xd': 60, 'first_choice': True}。copy.deepcopy 后 pop 出 'class',剩下的 key-value 作为构造参数。这是 ABU 全书统一的"dict 定义因子"模式(第 3 章买入因子、第 4 章卖出因子都是这样)。deepcopy 保证多进程下每个进程的 dict 独立(避免共享引用导致的并发问题)。
first_choice:True 的因子进 first_stock_pickers,其余进 stock_pickers。这两类因子的执行逻辑完全不同(下面 fit 详讲)。类型检测(L69)保证因子必须继承 AbuPickStockBase。
89 def _first_batch_fit(): 99 inner_first_choice_symbols = self.choice_symbols 101 with AbuMulPidProgress(len(self.first_stock_pickers), 'pick first_choice stocks complete') as progress: 102 for epoch, first_choice in enumerate(self.first_stock_pickers): 105 inner_first_choice_symbols = first_choice.fit_first_choice(self, inner_first_choice_symbols) 106 return inner_first_choice_symbols 108 def _batch_fit(): 118 with AbuMulPidProgress(len(self.choice_symbols), 'pick stocks complete') as progress: 120 inner_choice_symbols = [] 121 for epoch, target_symbol in enumerate(self.choice_symbols): 124 add = True 125 for picker in self.stock_pickers: 126 kl_pd = self.kl_pd_manager.get_pick_stock_kl_pd(target_symbol, picker.xd, picker.min_xd) 131 sub_add = picker.fit_pick(kl_pd, target_symbol) 132 if sub_add is False: 133 add = False 134 break # 一票否决,任一反对即 break 136 if add: 137 inner_choice_symbols.append(target_symbol) 138 return inner_choice_symbols 141 self.choice_symbols = _first_batch_fit() # 先首选层层过滤 143 self.choice_symbols = _batch_fit() # 再一票否决
两阶段:
阶段一:_first_batch_fit(首选,层层过滤)
fit_first_choice。inner_first_choice_symbols,下一个因子再对这个子集筛。fit_first_choice(批量高效,一次性处理整个池子),与普通的 fit_pick(逐只判断)不同。阶段二:_batch_fit(普通,一票否决)
fit_pick(kl_pd, symbol) 投票。两阶段的顺序:先首选(快速缩小池子),再一票否决(精细筛选)。如果首选阶段没有因子(first_stock_pickers 为空),_first_batch_fit 直接返回原 choice_symbols(L95-97);如果普通阶段没有因子,_batch_fit 也直接返回(L114-116)。这种"无因子即全通过"的退化,保证 fit 永远不会卡住。
💡 解剖要点:一票否决制体现了 ABU 选股的保守哲学——宁愿漏掉好股票(假阴性),也不要放进坏股票(假阳性)。因为坏股票会浪费回测资源和信号,而漏掉的好股票还有择时阶段的买入因子把关。首选制则解决"批量初筛"的效率问题——如果 1000 只全走一票否决(逐只取 KL + 多因子投票),太慢;先用首选因子的批量接口筛到几十只,再走一票否决,快得多。
下一节讲择时调度——AbuPickTimeMaster 怎么并行 + AbuPickTimeWorker 怎么用 kl_pd.apply(_task_loop) 时间驱动逐日推进,以及用日期差>60 识别月末的巧妙设计。