第 10 章 · 02 择时调度:时间驱动与自然周月 本节摘要:本节解剖 ABU 回测流水线的"择时调度中枢"—— + + ( 约 430 行择时部分)。选股筛出股票池后,择时对每只股票"逐日决策买卖"。 是并行总指挥:主进程 一次性拉满所有 KL → 切分 → 并行分发(子进程 不融合资金)→ 合并 orders/action → 主进程统一 (注意 L71 有个 应为 的 bug)。 是择时执行者: 合并回测前 1 年数据(供因子生成特征) → 预筛选周月因子。fit 是时间驱动的核心: 按交易日逐行推进。自然周月识别用巧妙技巧—— (每周五)、 (日期差>60 识别月末,不用时间 API 效率高)。
本节摘要:本节解剖 ABU 回测流水线的"择时调度中枢"——
AbuPickTimeMaster+AbuPickTimeWorker+ABuPickTimeExecute(AlphaBu约 430 行择时部分)。选股筛出股票池后,择时对每只股票"逐日决策买卖"。AbuPickTimeMaster.do_symbols_with_same_factors_process是并行总指挥:主进程batch_get_pick_time_kl_pd一次性拉满所有 KL →split_k_market切分 → 并行分发(子进程apply_capital=False不融合资金)→pd.concat合并 orders/action → 主进程统一apply_action_to_capital(注意 L71 有个==应为=的 bug)。AbuPickTimeWorker是择时执行者:combine_kl_pd合并回测前 1 年数据(供因子生成特征) →filter_long_task_factors预筛选周月因子。fit 是时间驱动的核心:kl_pd.apply(_task_loop, axis=1)按交易日逐行推进。自然周月识别用巧妙技巧——week_task=where(date_week==4,1,0)(每周五)、month_task=where(shift(-1).date-date>60,1,0)(日期差>60 识别月末,不用时间 API 效率高)。_day_task日任务主轴执行顺序:先卖后买(_task_attached_sell专属卖出 → sell_factors → buy_factors),避免高频日交易。本节会精读这套时间驱动引擎。
内容来源:原项目源码
abupy/AlphaBu/ABuPickTimeMaster.py(104 行)+ABuPickTimeWorker.py(324 行)+ABuPickTimeExecute.py(190 行)。
⚠️ 注意:本节重点理解三个精妙设计。一是"主进程拉满 KL + 子进程不融合资金 + 主进程统一融合"——避免多进程同时操作资金对象(线程不安全),资金融合只能在单进程做。二是自然周月的"日期差>60"识别——不用 datetime API,直接
20140801 - 20140731 = 70 > 60判断月末,运行效率比时间 API 快一个数量级。三是"先卖后买"——回测始终非高频(不当日买卖),卖出先于买入执行,避免同一天又买又卖造成假高频。另外 L71 有个源码 bug:==应为=,导致 FORCE_LOCAL 设置失效(但不影响正确性,因为子进程只读不写)。
阅读完本节,你应当能够:
do_symbols_with_same_factors_process:主进程 batch_get_pick_time_kl_pd 拉满 → split_k_market 切分 → 并行 apply_capital=False → pd.concat 合并 → 主进程 apply_action_to_capital。apply_capital=False(多进程同时操作资金对象不安全),以及为什么资金融合必须在主进程单进程做。== 应为 =),以及它为什么不影响正确性。combine_kl_pd 的作用(合并回测前 1 年数据,供因子生成 MA/ATR 等特征)和 filter_long_task_factors 的预筛选(初始化时筛周月因子,避免时间循环里反复 hasattr)。fit:g_natural_long_task 下添加 week_task/month_task 列 → kl_pd.apply(_task_loop, axis=1) 逐日推进。week_task=where(date_week==4,1,0)(每周五)/ month_task=where(shift(-1).date-date>60,1,0)(日期差>60,不用时间 API)。_task_loop 的执行顺序(月→周→日),以及 _day_task 的"先卖后买"(_task_attached_sell → sell_factors → buy_factors)。do_symbols_with_same_factors_process(ABuPickTimeMaster.py:29-103)是对外接口:
49 if kl_pd_manager is None: 50 kl_pd_manager = AbuKLManager(benchmark, capital) 52 kl_pd_manager.batch_get_pick_time_kl_pd(target_symbols, n_process=n_process_kl, show_progress=show_progress) 59 process_symbols = split_k_market(n_process_pick_time, market_symbols=target_symbols) 67 tmp_fetch_mode = ABuEnv.g_data_fetch_mode 68 if n_process_pick_time > 1: 71 ABuEnv.g_data_fetch_mode == EMarketDataFetchMode.E_DATA_FETCH_FORCE_LOCAL # bug: == 应为 = 74 p_nev = AbuEnvProcess() 76 out = parallel(delayed(do_symbols_with_same_factors)(choice_symbols, benchmark, buy_factors, sell_factors, 77 capital, apply_capital=False, 78 kl_pd_manager=kl_pd_manager, env=p_nev, ...) for choice_symbols in process_symbols) 82 ABuEnv.g_data_fetch_mode = tmp_fetch_mode 86 for sub_out in out: 88 sub_orders_pd, sub_action_pd, sub_all_fit_symbols_cnt = sub_out 89 orders_pd = pd.concat([orders_pd, sub_orders_pd]) # 合并 orders 90 action_pd = pd.concat([action_pd, sub_action_pd]) # 合并 action 96 action_pd = action_pd.sort_values(['Date', 'action']) # 按时间+行为排序 101 ABuTradeExecute.apply_action_to_capital(capital, action_pd, kl_pd_manager, ...) # 主进程融合资金
batch_get_pick_time_kl_pd 在主进程一次性把所有 target_symbols 的 KL 数据拉满,塞进 kl_pd_manager。这一步用 n_process_kl 个进程并行(数据采集阶段允许并行,因为读 HDF5 安全)。预拉满后,后续子进程直接从 manager 取,不用再联网/读盘。
apply_capital=False 是关键——子进程只生成 orders_pd 和 action_pd(交易行为记录),不把行为作用到资金对象上。原因:capital 是 AbuCapital 实例,内部维护资金时序,多进程同时操作会造成数据竞争(资金被并发改乱)。所以资金融合(apply_action_to_capital)必须在单进程做。
71 ABuEnv.g_data_fetch_mode == EMarketDataFetchMode.E_DATA_FETCH_FORCE_LOCAL
这行用 ==(等于判断)而不是 =(赋值)。== 是个表达式,返回布尔值但什么都不改——所以 g_data_fetch_mode 没被设成 FORCE_LOCAL。作者的意图是"并行时强制本地模式"(避免子进程联网),但 bug 导致设置失效。
为什么不影响正确性?因为主进程已经预拉满 KL(L52),子进程从 kl_pd_manager 取数据时走的是内存,不会触发联网。所以即使 g_data_fetch_mode 没改,实际行为也正确。但这是个真实的代码 bug,读源码时要警惕。L82 还有个 ABuEnv.g_data_fetch_mode = tmp_fetch_mode(恢复原值),但因为 L71 没真正改,这行恢复也没意义。
所有子进程完成后,合并 orders_pd 和 action_pd,按 Date+action 排序,最后在主进程调 apply_action_to_capital 把所有交易行为作用到 capital 上,生成资金时序。这是第 5 章讲过的 AbuCapital 核心流程。
AbuPickTimeWorker.__init__(ABuPickTimeWorker.py:37-64):
37 def __init__(self, cap, kl_pd, benchmark, buy_factors, sell_factors): 45 self.capital = cap 47 self.kl_pd = kl_pd 49 self.combine_kl_pd = ABuSymbolPd.combine_pre_kl_pd(self.kl_pd, n_folds=1) # 合并回测前 1 年 56 self.init_buy_factors(buy_factors) 58 self.init_sell_factors(sell_factors) 60 self.filter_long_task_factors() # 预筛选周月因子 62 self.orders = list() # 最终订单列表 64 self.task_pg = None # 进度条
combine_pre_kl_pd(kl_pd, n_folds=1) 把回测开始前 1 年的 KL 数据拼到 kl_pd 前面,形成 combine_kl_pd。为什么?因为买入因子要算 MA60、ATR21 等指标,这些指标需要历史数据预热。如果回测从 2016-01-01 开始,第一天的 MA60 算不出来(没有前 60 天数据)。combine_kl_pd 提供前 1 年数据让因子能算指标,但回测本身只针对 kl_pd 部分(combine_kl_pd 不参与交易)。
源码注释 L51-52 提到优化:如果不在乎效率,可以只在 g_enable_ml_feature 模式下开启 combine_kl_pd(因为 ML 特征需要更多历史)。默认总是开启。
308 def filter_long_task_factors(self): 315 self.week_buy_factors = list(filter(lambda buy_factor: hasattr(buy_factor, 'fit_week'), self.buy_factors)) 317 self.month_buy_factors = list(filter(lambda buy_factor: hasattr(buy_factor, 'fit_month'), self.buy_factors)) 320 self.week_sell_factors = list(filter(lambda sell_factor: hasattr(sell_factor, 'fit_week'), self.sell_factors)) 322 self.month_sell_factors = list(filter(lambda sell_factor: hasattr(sell_factor, 'fit_month'), self.sell_factors))
在初始化时一次性筛选出"支持周任务/月任务"的因子。为什么?因为时间驱动循环里每个交易日都要判断"这个因子支不支持 fit_week/fit_month",如果每次都 hasattr,几万个交易日 × 几个因子 × hasattr 的开销不小。预先在 init 时筛好,循环里直接迭代筛好的列表,效率高。这是 ABU 对热路径的优化。
fit(ABuPickTimeWorker.py:208-245)是择时的核心引擎:
208 def fit(self, *args, **kwargs): 212 if g_natural_long_task: 215 self.kl_pd['week_task'] = np.where(self.kl_pd.date_week == 4, 1, 0) 240 self.kl_pd['month_task'] = np.where(self.kl_pd.shift(-1)['date'] - self.kl_pd['date'] > 60, 1, 0) 242 self.kl_pd.apply(self._task_loop, axis=1)
self.kl_pd['week_task'] = np.where(self.kl_pd.date_week == 4, 1, 0)
date_week 是 0-6(周一到周日),== 4 即周五。每周五标记 week_task=1。为什么用周五而不是用 day % 5 == 0?因为后者是"每 5 个交易日",遇到节假日会错位;周五是"自然周的最后一个交易日"(周末不开盘),即使周中有节假日,周五依然是周末前最后一天。这就是"自然周"的含义——对齐日历周,不是交易日计数。
self.kl_pd['month_task'] = np.where(self.kl_pd.shift(-1)['date'] - self.kl_pd['date'] > 60, 1, 0)
源码注释(L217-239)详细解释:
2014-07-28 1.0 2014-07-29 1.0 2014-07-30 1.0 2014-07-31 70.0 ← 月末!下一天是 2014-08-01,20140801-20140731=70 2014-08-01 3.0 2014-08-04 1.0
shift(-1) 是"下一天",下一天.date - 今天.date 得到日期差(整数相减)。月末时,今天 7/31 明天 8/1,20140801 - 20140731 = 70 > 60,标记 month_task=1。
**为什么用 60 而不是精确判断?**因为这是整数相减(20140801 - 20140731),不是 datetime 运算。月末到下月初的差值(70)远大于普通日与日的差值(1-3,跨周末最多 3)。60 是个安全阈值——普通日差最多 3(周五到周一),月末差至少 70(31 号到 1 号的整数差),60 居中,完美区分。
**为什么不用 datetime API?**源码注释明说"没有使用时间 api,因为这样做运行效率快"。datetime 运算(解析、比较)比整数减法慢一个数量级。对几百只股票 × 几百个交易日的批量回测,这个优化省下的时间可观。这是 ABU 对性能的极致追求。
这是时间驱动的核心——apply(axis=1) 按行(按交易日)迭代,每行的 today 是一个 Series(那天的 OHLCV + week_task + month_task)。_task_loop 对每个 today 执行任务。这种"逐日推进"的引擎,忠实模拟了真实交易——你只能在今天的信息下决策,不能用未来数据。
_task_loop(ABuPickTimeWorker.py:172-205):
172 def _task_loop(self, today): 185 day_cnt = today.key 187 today.exec_week = today.week_task == 1 if g_natural_long_task else day_cnt % 5 == 0 189 today.exec_month = today.month_task == 1 if g_natural_long_task else day_cnt % 20 == 0 191 if day_cnt == 0 and not today.exec_week: 193 self._task_attached_ps(today, is_week=True) # 第一天初始化周选股池 194 if day_cnt == 0 and not today.exec_month: 196 self._task_attached_ps(today, is_week=False) # 第一天初始化月选股池 198 if today.exec_month: 200 self._month_task(today) # 月任务 201 if today.exec_week: 203 self._week_task(today) # 周任务 205 self._day_task(today) # 日任务(每天必跑)
执行顺序:月 → 周 → 日。
g_natural_long_task=True(默认)用自然周月(week_task/month_task 列),False 用计数(每 5 天/20 天)。周月任务里不建议生成买单和卖单(注释 L87、L107),执行单全部在日任务完成。周月任务主要做"选股池更新"和"因子状态刷新",通过 today.exec_week/today.exec_month 让日任务知道"今天是周/月边界"。
_day_task(ABuPickTimeWorker.py:120-141)是日任务的主轴:
120 def _day_task(self, today): 127 self._task_attached_sell(today, how='day') # 1. 专属卖出(买入因子自带的卖出) 130 for sell_factor in self.sell_factors: 132 sell_factor.read_fit_day(today, self.orders) # 2. 卖出因子(先于买入) 135 for buy_factor in self.buy_factors: 137 if not buy_factor.lock_factor: 139 order = buy_factor.read_fit_day(today) # 3. 买入因子(后于卖出) 140 if order and order.order_deal: 141 self.orders.append(order)
执行顺序:专属卖出 → 卖出因子 → 买入因子。
买入因子可以"自带"卖出因子(比如"突破买入 + 这个突破失败就卖")。_task_attached_sell 对每个买入因子,找出它自带的 sell_factors,针对该因子的历史订单执行卖出。关键:即使买入因子被锁(lock_factor=True),它的专属卖出不能锁——被锁只是不能开新仓,该平的旧仓还是要平。
注释 L129:"注意回测模式下始终非高频,非当日买卖,不区分美股,A股市场,卖出因子要先于买入因子的执行"。
为什么先卖后买?因为回测非高频(不当日买卖)。如果先买后卖,可能出现"早上买入下午卖出"的假高频日交易,这在真实交易里受 T+1(A 股)或资金限制做不到。先卖后买保证:今天的卖出是针对昨天及以前的持仓,今天的买入是新建仓,两者不冲突。这是 ABU 对真实交易约束的模拟。
买入因子有 lock_factor 属性。被锁的买入因子不执行 fit_day(不发买单)。锁定由周/月任务的专属选股因子(_task_attached_ps)触发——如果选股因子判断"当前不该买这个因子",就锁住它。这是"选股"层面对"择时"的控制。
ABuPickTimeExecute.py:81-163 是多 symbol 的循环包装:
107 for epoch, target_symbol in enumerate(target_symbols): 117 kl_pd = kl_pd_manager.get_pick_time_kl_pd(target_symbol) 118 ret, fit_error = _do_pick_time_work(capital, p_buy_factors, p_sell_factors, kl_pd, benchmark, ...) 121 except Exception as e: 122 logging.exception(e) 123 continue # 异常容错,跳过该 symbol 125 if ret is None and back_target_symbols is not None: 127 if fit_error == EFitError.NO_ORDER_GEN: 129 r_all_fit_symbols_cnt += 1 # 无 order 也统计 130 while True: 134 target_symbol = back_target_symbols.pop() # 补位机制 140 if ret is not None: break
两个关键设计:
back_target_symbols 弹一个补上。这保证最终回测的 symbol 数量达标(比如要回测 100 只,有 3 只失败,从补位池补 3 只)。EFitError 枚举(L29-43)区分错误类型:NET_ERROR/DATE_ERROR/NO_ORDER_GEN/OTHER_ERROR。ABuEnv.g_data_fetch_mode == E_DATA_FETCH_FORCE_LOCAL 用 ==(判断)而非 =(赋值),设置失效。但不影响正确性——主进程已预拉满 KL,子进程走内存不联网。L82 的恢复也成空操作。下一节(全书最后一节)讲 ipywidgets GUI——ABU 怎么用 4716 行代码把回测/选股/因子/UMP/网格搜索/交叉验证全部做成 Tab 交互界面,以及全书的系统回顾。