feat(data): 恢复双源(task#79)—撮合raw+策略qfq, 分红除权准确
用户要模拟=回测准确: raw除权缺口致MA假信号, 必须双源。
- data_source: qfq→qfq_dir(干净qfq), raw→raw_dir; _check_adjust_cfg(cfg提供才校验)
- engine 双bar流: step(raw_bars,qfq_bars)撮合/盯市raw+策略on_bar qfq; run zip(raw,qfq)
- live_orchestrator: warmup用qfq(信号am); 去adjust参数(双源固定)
- raw_redownload --adjust(''raw/'qfq'); config qfq_dir
- 113/113通过
This commit is contained in:
+27
-23
@@ -1,10 +1,9 @@
|
||||
"""PaperEngine 模拟盘主循环(逐根 bar 重放 + 双层记账 + 持久化,spec §4/§9)。
|
||||
|
||||
run():逐 bar → T+1 解冻 → 撮合上一根 next_open pending(用当前 bar)→ 喂策略
|
||||
on_bar 收新单 → current_close 当根撮合 / next_open 缓冲到下根 → 盯市 → 入库。
|
||||
|
||||
adjust 默认 raw(真实价):撮合/涨跌停/成交价/信号共用一套真实价 bar。
|
||||
除权缺口对 MA 信号的影响留分期项(分红除权)处理。
|
||||
双源(分红除权准确方案,task #79 恢复):
|
||||
- 撮合/涨跌停/盯市用 **raw**(真实价,涨跌停/成交真实)
|
||||
- 策略 on_bar 信号用 **qfq**(前复权,无除权缺口 → MA 信号准)
|
||||
run() 双迭代器 zip(raw, qfq) 同日期对齐;step(raw_bars, qfq_bars)。
|
||||
"""
|
||||
import logging
|
||||
|
||||
@@ -41,7 +40,7 @@ class PaperEngine:
|
||||
def __init__(self, account: Account, runners: list[StrategyRunner],
|
||||
data_source, cfg, db_path: str, account_id: int,
|
||||
symbols: list[str], start: str, end: str,
|
||||
interval: str = "d", adjust: str = "raw") -> None:
|
||||
interval: str = "d") -> None:
|
||||
self.account = account
|
||||
self.runners = runners
|
||||
self.data_source = data_source
|
||||
@@ -52,34 +51,34 @@ class PaperEngine:
|
||||
self.start = start
|
||||
self.end = end
|
||||
self.interval = interval
|
||||
self.adjust = adjust
|
||||
|
||||
def step(self, bar_date, bars, prev_close, pending):
|
||||
"""单根 bar 推进(回放 run 循环调;C-S3 实走 scheduler 每日调)。
|
||||
def step(self, bar_date, raw_bars, qfq_bars, prev_close, pending):
|
||||
"""单根 bar 推进(回放 run 循环调;实走 live_step 调)。
|
||||
|
||||
返回 (新 pending, 当根 closes)——实走每日喂当日 bar 调一次。
|
||||
撮合/盯市用 raw_bars(真实价);策略 on_bar 用 qfq_bars(信号准)。
|
||||
返回 (新 pending, 当根 closes)。
|
||||
"""
|
||||
self._bar_count = getattr(self, "_bar_count", 0) + 1
|
||||
self.account.unfreeze_all()
|
||||
for r in self.runners:
|
||||
r.unfreeze_all()
|
||||
# 1. 撮合上一根 pending(next_open,用当前 bar)
|
||||
# 1. 撮合上一根 pending(next_open,用当日 raw bar)
|
||||
if pending:
|
||||
for order, runner in pending:
|
||||
self._match(order, runner, bars, prev_close, bar_date)
|
||||
self._match(order, runner, raw_bars, prev_close, bar_date)
|
||||
pending = []
|
||||
# 2. 喂策略 on_bar → 收新单
|
||||
# 2. 喂策略 on_bar(qfq 信号)→ 收新单 → 当根撮合 raw / 缓冲 next_open
|
||||
for runner in self.runners:
|
||||
sym = runner.symbol
|
||||
if sym and sym in bars:
|
||||
runner.paper_cta_engine.on_bar(bars[sym])
|
||||
if sym and sym in qfq_bars:
|
||||
runner.paper_cta_engine.on_bar(qfq_bars[sym])
|
||||
for order in runner.paper_cta_engine.pop_orders():
|
||||
if order.match_session == MatchSession.NEXT_OPEN:
|
||||
pending.append((order, runner))
|
||||
else: # current_close 当根撮合
|
||||
self._match(order, runner, bars, prev_close, bar_date)
|
||||
# 3. 盯市 + 入库
|
||||
closes = {s: bars[s].close_price for s in bars}
|
||||
else: # current_close 当根撮合(raw)
|
||||
self._match(order, runner, raw_bars, prev_close, bar_date)
|
||||
# 3. 盯市 raw + 入库
|
||||
closes = {s: raw_bars[s].close_price for s in raw_bars}
|
||||
self.account.mark_to_market(closes)
|
||||
save_daily_balance(
|
||||
self.db_path, self.account_id, str(bar_date),
|
||||
@@ -90,12 +89,17 @@ class PaperEngine:
|
||||
return pending, closes
|
||||
|
||||
def run(self) -> None:
|
||||
"""双源 zip(raw, qfq) 同日期对齐,逐根 step。"""
|
||||
prev_close: dict[str, float] = {}
|
||||
pending: list = [] # [(order, runner)] next_open 待下根撮合
|
||||
for bar_date, bars in self.data_source.iter_bars(
|
||||
self.symbols, self.start, self.end, self.interval, self.adjust, None
|
||||
):
|
||||
pending, closes = self.step(bar_date, bars, prev_close, pending)
|
||||
raw_iter = self.data_source.iter_bars(
|
||||
self.symbols, self.start, self.end, self.interval, "raw", None
|
||||
)
|
||||
qfq_iter = self.data_source.iter_bars(
|
||||
self.symbols, self.start, self.end, self.interval, "qfq", None
|
||||
)
|
||||
for (rdate, raw_bars), (_qdate, qfq_bars) in zip(raw_iter, qfq_iter):
|
||||
pending, closes = self.step(rdate, raw_bars, qfq_bars, prev_close, pending)
|
||||
prev_close = closes
|
||||
|
||||
def _match(self, order, runner, bars, prev_close, bar_date) -> None:
|
||||
|
||||
Reference in New Issue
Block a user