fix(pipeline): 冒烟实证修——桥扁DataFrame行走+env全链配方+epoch退化守卫 [vps]
CI/CD / test (push) Successful in 26s
CI/CD / nas-deploy (push) Successful in 9s
CI/CD / nas-verify (push) Successful in 8s

三案皆 VPS 真桥只读探针实证(09-26 冒烟):
- is_price_source_vps:minute_rows_to_map 形态归一(扁 DataFrame[index,close]
  实测形态为准,兼容遗留嵌套 dict);_minute_closes 改逐票调桥(扁形态无
  code 列,批量切片无法按票归属);_ensure_bridge_env 补 DEFAULT_DATA_
  PROVIDER=miniqmt + PYTHONPATH 全链(bridge/vnpy 根/vnpy_v4.4.0/vnpy_qmt_
  v0.3.3,SANGUO_VNPY_ROOT 可覆)——缺任一 xtdata import 即断。
- order_log_parser._t_raw_ts:epoch sanity floor=2000-01-01(生产实见
  t_raw=0 退化单 3565 条→1970 幽灵日;退化/负/垃圾回退行首时间戳),
  毫秒分支移入 floor 内评估。
- 测试 +3(扁/嵌套形态×2、epoch 退化回退×1)。

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
2026-09-26 14:07:22 +08:00
parent abbee95994
commit b5aa18cf8d
4 changed files with 89 additions and 28 deletions
+6 -4
View File
@@ -38,13 +38,15 @@ TERMINAL_MAP = {"filled": "fill", "rejected": "reject", "failed": "reject",
def _t_raw_ts(t_raw: str, line_ts: str) -> str:
"""柜台 t_raw(epoch)→时刻;fromtimestamp 走机器本地时区——本舰队全 UTC+8
成立,跨时区复用需显式 tz(infra 注记①)。"""
成立,跨时区复用需显式 tz(infra 注记①)。sanity floor=2000-01-01:
生产实见 t_raw=0 退化单(日志 3565 条)→1970 幽灵日,回退行首时间戳。"""
if t_raw not in ("NONE", ""):
try:
v = float(t_raw)
if v > 1e12: # 毫秒形态
v /= 1000.0
return datetime.fromtimestamp(v).strftime("%Y-%m-%d %H:%M:%S")
if v > 946684800: # sanity floor(2000 后才可信;退化 0/负/垃圾回退)
if v > 1e12: # 毫秒形态(floor 内评估,~1.7e12 双过)
v /= 1000.0
return datetime.fromtimestamp(v).strftime("%Y-%m-%d %H:%M:%S")
except (ValueError, OSError, OverflowError):
pass
return line_ts
+57 -24
View File
@@ -1,8 +1,9 @@
# -*- coding: utf-8 -*-
"""VPS 价格源适配器(偏差日报三价链锚点;Mac 不可跑——桥/dbbardata 仅 VPS)。
分钟 close:桥 xtdata shim(PYTHONPATH 前置 C:\\sanguo_bigqmt\\xtquant_bridge +
env BIGQMT_ACCOUNT_ID/BIGQMT_REDIS_PASSWORD——口令运行时读 C:\\redis\\redis.conf,
分钟 close:桥 xtdata shim(PYTHONPATH 全链 bridge+vnpy 根+vnpy_v4.4.0+
vnpy_qmt_v0.3.3,env DEFAULT_DATA_PROVIDER/BIGQMT_ACCOUNT_ID/
BIGQMT_REDIS_PASSWORD——口令运行时读 C:\\redis\\redis.conf,
永不落盘/回显)。日线 OHLC/prev_close:quant_trading.db dbbardata(datareader
同款列名)。symbol 规范化:日志/台账四种形态 → xtdata 的 XXXXXX.SH/SZ。
"""
@@ -12,9 +13,11 @@ import os
import re
import sqlite3
from datetime import datetime, timedelta
from typing import Any
_BRIDGE_PYROOT = os.environ.get(
"SANGUO_BRIDGE_PYROOT", r"C:\sanguo_bigqmt\xtquant_bridge")
_VNPY_ROOT = os.environ.get("SANGUO_VNPY_ROOT", r"C:\sanguo_vnpy_v2")
_ACCOUNT = os.environ.get("BIGQMT_ACCOUNT_ID", "66639661")
_DAILY_DB = os.environ.get(
"SANGUO_DAILY_DB", r"C:\sanguo_vnpy_v2\data\quant_trading.db")
@@ -42,43 +45,73 @@ def _read_redis_password() -> str:
def _ensure_bridge_env() -> None:
"""桥 env 配方(冒烟 09-26 实证):DEFAULT_DATA_PROVIDER=miniqmt +
PYTHONPATH 全链(bridge/vnpy 根/vnpy_v4.4.0/vnpy_qmt_v0.3.3)——缺任一项
xtdata import 即报 无法连接xtquant服务。逐项 append,已有项不重复。"""
parts = [p for p in os.environ.get("PYTHONPATH", "").split(os.pathsep)
if p]
if _BRIDGE_PYROOT not in parts:
os.environ["PYTHONPATH"] = os.pathsep.join([_BRIDGE_PYROOT] + parts)
for p in (_BRIDGE_PYROOT, _VNPY_ROOT,
os.path.join(_VNPY_ROOT, "vnpy_v4.4.0"),
os.path.join(_VNPY_ROOT, "vnpy_qmt_v0.3.3")):
if p not in parts:
parts.append(p)
os.environ["PYTHONPATH"] = os.pathsep.join(parts)
os.environ.setdefault("DEFAULT_DATA_PROVIDER", "miniqmt")
os.environ.setdefault("BIGQMT_ACCOUNT_ID", _ACCOUNT)
if "BIGQMT_REDIS_PASSWORD" not in os.environ:
os.environ["BIGQMT_REDIS_PASSWORD"] = _read_redis_password()
def minute_rows_to_map(obj: Any, code: str) -> dict[tuple[str, str], float]:
"""桥 get_market_data 返回 → {(code,'HH:MM'): close}。
实测返回=扁 DataFrame[index, close](冒烟 09-26 实证;时间在 index 列
或行索引皆出值);兼容遗留嵌套 dict {code:{ts:close}} 形态。分钟键=
digits 解析("20260924 14:59:00" 与 "2026-09-24 14:59:00" 皆出 "14:59")。
"""
out: dict[tuple[str, str], float] = {}
if hasattr(obj, "iterrows"): # 扁 DataFrame
for idx, row in obj.iterrows():
v = row.get("close") if hasattr(row, "get") else row["close"]
if v is None or v != v:
continue
ts_src = row.get("index", idx) if hasattr(row, "get") else idx
digits = re.sub(r"\D", "", str(ts_src))
if len(digits) >= 12:
minute = f"{digits[8:10]}:{digits[10:12]}"
else:
continue
out[(code, minute)] = float(v)
return out
sub = obj.get(code) if hasattr(obj, "get") else None # 遗留嵌套形态
if sub is None:
return out
for idx, v in (sub.items() if hasattr(sub, "items") else []):
if v is None or v != v:
continue
digits = re.sub(r"\D", "", str(idx))
if len(digits) >= 12:
out[(code, f"{digits[8:10]}:{digits[10:12]}")] = float(v)
return out
def _minute_closes(symbols: list[str], day: str) -> dict[tuple[str, str], float]:
"""桥拉一日分钟 close → {(symbol, 'HH:MM'): close}(当日锚点一次拉齐)。"""
"""桥拉一日分钟 close → {(symbol, 'HH:MM'): close}(当日锚点一次拉齐)。
逐票调桥(冒烟 09-26 实证:返回扁 DataFrame 无 code 列,批量切片无法
按票归属);形态解析归一走 minute_rows_to_map。
"""
_ensure_bridge_env()
from xtquant import xtdata # noqa: WPS433(shim 需 env 先就位,延迟导入)
start = f"{day} 09:00:00"
end = f"{day} 15:30:00"
out: dict[tuple[str, str], float] = {}
codes = sorted({normalize_symbol(s) for s in symbols})
for i in range(0, len(codes), 50):
df = xtdata.get_market_data(
field_list=["close"], stock_list=codes[i:i + 50], period="1m",
for code in codes:
result = xtdata.get_market_data(
field_list=["close"], stock_list=[code], period="1m",
start_time=start, end_time=end)
if df is None or len(df) == 0:
continue
closes = df["close"] if hasattr(df, "keys") and "close" in df else df
for code in codes[i:i + 50]:
try:
sub = closes[code]
except (KeyError, TypeError, IndexError):
continue
for idx, v in sub.items() if hasattr(sub, "items") else []:
digits = re.sub(r"\D", "", str(idx))
if len(digits) >= 12: # "2026-09-18 09:35:00"/"20260918093500"
minute = f"{digits[8:10]}:{digits[10:12]}"
else: # 紧凑短形态兜底 "…0935"
minute = str(idx)[-5:-3] + ":" + str(idx)[-2:]
if v == v: # NaN 剔除
out[(code, minute)] = float(v)
out.update(minute_rows_to_map(result, code))
return out
@@ -12,3 +12,23 @@ def test_normalize_symbol_variants():
assert nz("000661.SZ") == "000661.SZ"
assert nz("sz000661") == "000661.SZ"
assert nz("000661.XSHE") == "000661.SZ"
def test_minute_rows_to_map_flat_dataframe():
import pandas as pd
from scripts.pipeline.is_price_source_vps import minute_rows_to_map
df = pd.DataFrame({"index": ["20260924 14:59:00", "20260924 15:00:00",
"2026-09-24 09:31:00"],
"close": [44.90, 44.86, 44.75]})
m = minute_rows_to_map(df, "600276.SH")
assert m == {("600276.SH", "14:59"): 44.90,
("600276.SH", "15:00"): 44.86,
("600276.SH", "09:31"): 44.75}
def test_minute_rows_to_map_legacy_nested_and_nan():
from scripts.pipeline.is_price_source_vps import minute_rows_to_map
legacy = {"600276.SH": {"2026-09-24 10:00:00": 44.8,
"bad": None}}
m = minute_rows_to_map(legacy, "600276.SH")
assert m == {("600276.SH", "10:00"): 44.8}
+6
View File
@@ -63,3 +63,9 @@ def test_classify_reject_two_taxonomy():
assert f("sell", 10.0, 9.5, 11.0, 9.0, bought_today=False) == "TIMEOUT"
assert f("buy", 11.0, 10.0, 11.0, 9.0, bought_today=False) == "LIMIT_UP"
assert f("buy", 10.5, 10.0, 11.0, 9.0, bought_today=False) == "CASH"
def test_t_raw_degenerate_epoch_falls_back():
from sanguo_portfolio.order_log_parser import _t_raw_ts
assert _t_raw_ts("0", "2026-09-18 09:35:21") == "2026-09-18 09:35:21"
assert _t_raw_ts("1789695321", "2026-09-18 09:35:21").startswith("2026-09-18")