perf(data): TickFlow 取数切换 as_dataframe=False 列式直转; 修复除权因子日期 UTC 偏移一天

- 全部 K 线取数点改走 CompactKlineData 列式直转 polars (无 pandas 中转,
  全市场单轮本地转换 ~1.1s → ~0.1s), 时区口径统一 _timestamp_to_beijing_datetime
- 除权因子 trade_date 原取 UTC 日期, 北京零点事件整体早一天 → 转北京墙钟后取日期
- 盘中分钟增量间隔上限 300s → 120s (universe 仅回最新 3 根, 超限必留缺口)
This commit is contained in:
shy3130
2026-08-30 22:29:33 +08:00
parent 1a7e6ad32c
commit d29f9c0ea8
9 changed files with 139 additions and 88 deletions
+8 -1
View File
@@ -71,8 +71,15 @@ def normalize_adj_factors(data, source: str = "tickflow") -> pl.DataFrame: # no
df = df.rename({k: v for k, v in rename_map.items() if k in df.columns})
if "trade_date" in df.columns:
if df.schema["trade_date"] in {pl.Int64, pl.Int32, pl.UInt64, pl.UInt32, pl.Float64, pl.Float32}:
# 毫秒时间戳 → 北京墙钟日期 (直接 from_epoch().dt.date() 是 UTC 日期,
# 除权事件时间戳为北京零点 = UTC 前一日 16:00, 会整体早一天)。
df = df.with_columns(
pl.from_epoch(pl.col("trade_date").cast(pl.Int64), time_unit="ms").dt.date().alias("trade_date")
pl.from_epoch(pl.col("trade_date").cast(pl.Int64), time_unit="ms")
.dt.replace_time_zone("UTC")
.dt.convert_time_zone("Asia/Shanghai")
.dt.replace_time_zone(None)
.dt.date()
.alias("trade_date")
)
else:
df = df.with_columns(pl.col("trade_date").cast(pl.Date, strict=False))
+12 -13
View File
@@ -53,25 +53,24 @@ class TickFlowProvider:
"period": "1d",
"adjust": "none",
"count": 10000 if start_time and end_time else 250,
"as_dataframe": True,
"as_dataframe": False,
"show_progress": False,
}
if start_time and end_time:
from app.services.kline_sync import _datetime_to_ms
from app.services.kline_sync import _compact_klines_to_df, _datetime_to_ms, _timestamp_to_beijing_datetime
kwargs["start_time"] = _datetime_to_ms(start_time)
kwargs["end_time"] = _datetime_to_ms(end_time)
raw = tf.klines.batch(symbols, **kwargs)
frames: list[pl.DataFrame] = []
if isinstance(raw, dict):
for sym, sub in raw.items():
normalized = normalize_daily(sub, default_symbol=sym, source=self.name)
if not normalized.is_empty():
frames.append(normalized)
else:
normalized = normalize_daily(raw, source=self.name)
if not normalized.is_empty():
frames.append(normalized)
return pl.concat(frames, how="diagonal_relaxed") if frames else pl.DataFrame()
from app.services.kline_sync import _compact_klines_to_df, _timestamp_to_beijing_datetime
raw = tf.klines.batch(symbols, **kwargs)
# False 直转: 列数组→polars (无 pandas 中转), 加北京墙钟 datetime 列
# (normalize_daily 映射为 date); 保留 timestamp 原列 — normalize_daily
# 会把它改名为 quote_ts (盘后校验/量比折算用), 不能像 kline_sync 路径那样丢弃。
seg = _compact_klines_to_df(raw)
if seg.is_empty():
return pl.DataFrame()
seg = seg.with_columns(_timestamp_to_beijing_datetime(pl.col("timestamp")).alias("datetime"))
return normalize_daily(seg, source=self.name)
def get_adj_factors(
self,
+92 -53
View File
@@ -118,25 +118,24 @@ def sync_daily_batch(symbols: list[str],
start_time=_datetime_to_ms(start_time),
end_time=_datetime_to_ms(end_time),
count=10000,
as_dataframe=True, show_progress=False,
as_dataframe=False, show_progress=False,
)
else:
raw = tf.klines.batch(chunk, period="1d", count=count or 250, adjust="none",
as_dataframe=True, show_progress=False)
as_dataframe=False, show_progress=False)
except Exception as e: # noqa: BLE001
logger.warning("batch fetch failed for %d symbols (chunk %d/%d): %s",
len(chunk), i + 1, len(chunks), e)
failed_syms.extend(chunk)
continue
# 兼容两种形态:dict[sym → df] 和扁平 df
if isinstance(raw, dict):
for sym, sub in raw.items():
if sub is None or len(sub) == 0:
continue
out.append(_normalize_daily(sub, default_symbol=sym))
elif raw is not None and len(raw) > 0:
out.append(_normalize_daily(raw))
# False 直转: timestamp(UTC 毫秒) → 北京墙钟 datetime 列,
# _normalize_daily 将 datetime 映射为 date — 与 SDK True 路径的
# trade_date 字符串列同口径 (fromtimestamp(ts/1000, Asia/Shanghai))。
seg = _compact_klines_to_df(raw)
if not seg.is_empty():
seg = seg.with_columns(_timestamp_to_beijing_datetime(pl.col("timestamp")).alias("datetime")).drop("timestamp")
out.append(_normalize_daily(seg))
if on_chunk_done:
on_chunk_done(i + 1, len(chunks))
@@ -308,8 +307,11 @@ def _normalize_adj_factor(raw) -> pl.DataFrame:
df = df.rename(rename_map)
if "trade_date" in df.columns:
if df.schema["trade_date"] in {pl.Int64, pl.Int32, pl.UInt64, pl.UInt32, pl.Float64, pl.Float32}:
# 毫秒时间戳 → 北京墙钟日期。不能直接 from_epoch().dt.date():
# 那是 UTC 日期, 除权事件时间戳为北京零点 (= UTC 前一日 16:00),
# 会整体早一天 (与 SDK True 路径 fromtimestamp(ts/1000, Asia/Shanghai) 不一致)。
df = df.with_columns(
pl.from_epoch(pl.col("trade_date").cast(pl.Int64), time_unit="ms").dt.date().alias("trade_date")
_timestamp_to_beijing_datetime(pl.col("trade_date")).dt.date().alias("trade_date")
)
else:
df = df.with_columns(pl.col("trade_date").cast(pl.Date, strict=False))
@@ -377,8 +379,8 @@ def sync_adj_factor(symbols: list[str], repo: KlineRepository,
default_rpm_when_unset=False,
)
# 构建 SDK 参数
sdk_kwargs: dict = {"as_dataframe": True, "batch_size": limit.batch, "show_progress": False}
# 构建 SDK 参数 (False: _normalize_adj_factor 的 dict 分支原生支持)
sdk_kwargs: dict = {"as_dataframe": False, "batch_size": limit.batch, "show_progress": False}
if start_time:
sdk_kwargs["start_time"] = _datetime_to_ms(start_time)
if end_time:
@@ -579,6 +581,53 @@ def _normalize_minute(df_in, default_symbol: str | None = None) -> pl.DataFrame:
return df.select(keep)
def _compact_klines_to_df(raw, default_symbol: str | None = None) -> pl.DataFrame:
"""SDK as_dataframe=False 的 K 线原始数据 (CompactKlineData, 列式) → polars 直转。
raw 两种形态: {symbol: Compact} (batch/universe) 或顶层 Compact (单标的)。
列数组本身就是 polars 的原生形状 — 无 pandas 中转, 也无 SDK True 路径对
每根 K 做的 datetime.fromtimestamp + trade_date/trade_time 字符串拼接。
这里只搬运原始列 (timestamp 毫秒 Int64 + OHLCV), datetime/类型收口交给
_normalize_minute / _normalize_daily — 与 True 路径同一契约。
"""
if isinstance(raw, dict) and "timestamp" in raw:
items: list[tuple[str | None, dict]] = [(None, raw)]
elif isinstance(raw, dict):
items = list(raw.items())
else:
return pl.DataFrame()
frames: list[pl.DataFrame] = []
for sym, kd in items:
if not isinstance(kd, dict):
continue
ts = kd.get("timestamp") or []
if not ts:
continue
key = sym if sym is not None else default_symbol
cols: dict[str, object] = {"timestamp": ts}
if key:
cols["symbol"] = [key] * len(ts)
for field in ("open", "high", "low", "close", "volume", "amount"):
values = kd.get(field)
if values is not None:
cols[field] = values
frames.append(pl.DataFrame(cols))
return pl.concat(frames, how="diagonal_relaxed") if frames else pl.DataFrame()
def _timestamp_to_beijing_datetime(col: pl.Expr) -> pl.Expr:
"""timestamp (UTC 毫秒) → 北京墙钟 naive datetime 表达式。
与 SDK True 路径 trade_date 列同口径: datetime.fromtimestamp(ts/1000, Asia/Shanghai)。
"""
return (
pl.from_epoch(col.cast(pl.Int64), time_unit="ms")
.dt.replace_time_zone("UTC")
.dt.convert_time_zone("Asia/Shanghai")
.dt.replace_time_zone(None)
)
def _datetime_to_ms(dt: datetime) -> int:
"""datetime → 毫秒时间戳 (供 SDK start_time / end_time 使用)。"""
return int(dt.timestamp() * 1000)
@@ -786,23 +835,19 @@ def sync_minute_batch(
end_time=_datetime_to_ms(cur_end),
count=10000,
adjust="forward",
as_dataframe=True, show_progress=False,
as_dataframe=False, show_progress=False,
)
else:
raw = tf.klines.batch(chunk, period="1m", count=count or 1200,
adjust="forward",
as_dataframe=True, show_progress=False)
as_dataframe=False, show_progress=False)
except Exception as e: # noqa: BLE001
logger.warning("minute batch fetch failed for %d symbols: %s", len(chunk), e)
continue
if isinstance(raw, dict):
for sym, sub in raw.items():
if sub is None or len(sub) == 0:
continue
seg_out.append(_normalize_minute(sub, default_symbol=sym))
elif raw is not None and len(raw) > 0:
seg_out.append(_normalize_minute(raw))
seg = _normalize_minute(_compact_klines_to_df(raw))
if not seg.is_empty():
seg_out.append(seg)
if on_chunk_done:
on_chunk_done(step, total_steps, seg_label)
@@ -862,17 +907,6 @@ def intraday_monitor_support(capset: CapabilitySet | None) -> dict[str, object]:
}
def _normalize_intraday_raw(raw, default_symbol: str | None = None) -> list[pl.DataFrame]:
frames: list[pl.DataFrame] = []
if isinstance(raw, dict):
for symbol, sub in raw.items():
if sub is not None and len(sub) > 0:
frames.append(_normalize_minute(sub, default_symbol=str(symbol)))
elif raw is not None and len(raw) > 0:
frames.append(_normalize_minute(raw, default_symbol=default_symbol))
return [frame for frame in frames if not frame.is_empty()]
def fetch_intraday_monitor_batch(
symbols: list[str], capset: CapabilitySet | None, *, now: datetime | None = None,
) -> pl.DataFrame:
@@ -900,16 +934,20 @@ def fetch_intraday_monitor_batch(
if source == "intraday_batch":
limits = capset.limits(Cap.INTRADAY_BATCH) if capset else None
raw = tf.klines.intraday_batch(
symbols, count=300, as_dataframe=True, show_progress=False,
symbols, count=300, as_dataframe=False, show_progress=False,
batch_size=limits.batch if limits and limits.batch else 100,
)
frames.extend(_normalize_intraday_raw(raw))
df = _normalize_minute(_compact_klines_to_df(raw))
elif source == "intraday_single":
raw = tf.klines.intraday(symbols[0], count=300, as_dataframe=True)
frames.extend(_normalize_intraday_raw(raw, default_symbol=symbols[0]))
raw = tf.klines.intraday(symbols[0], count=300, as_dataframe=False)
df = _normalize_minute(_compact_klines_to_df(raw, default_symbol=symbols[0]))
elif source == "minute_single":
raw = tf.klines.get(symbols[0], period="1m", count=300, as_dataframe=True)
frames.extend(_normalize_intraday_raw(raw, default_symbol=symbols[0]))
raw = tf.klines.get(symbols[0], period="1m", count=300, as_dataframe=False)
df = _normalize_minute(_compact_klines_to_df(raw, default_symbol=symbols[0]))
else:
df = pl.DataFrame()
if not df.is_empty():
frames.append(df)
except Exception as e: # noqa: BLE001
logger.warning("intraday monitor fetch failed (%s, %d symbols): %s", source, len(symbols), e)
return pl.DataFrame()
@@ -954,10 +992,11 @@ def fetch_intraday_full_market_burst(
# 已成功的数据一起拖垮 (pool.map 迭代中抛异常会废弃全部已收 frames)
try:
raw = tf.klines.intraday_batch(
chunk, count=count, as_dataframe=True, show_progress=False,
chunk, count=count, as_dataframe=False, show_progress=False,
batch_size=len(chunk),
)
return (_normalize_intraday_raw(raw), None)
seg = _normalize_minute(_compact_klines_to_df(raw))
return ([seg] if not seg.is_empty() else [], None)
except Exception as e:
return ([], e)
@@ -1007,17 +1046,22 @@ def fetch_intraday_universe_increment(
的 intraday.batch 脉冲 (请求量 28→1, 传输量 ~40 倍降)。缺口回补
(冷启动/长时间断档/全天修复) 仍走 fetch_intraday_full_market_burst。
返回 (增量分钟K, 请求数); 拉取失败返回空 df 由调用方按失败轮处理。
as_dataframe=False: raw 的 CompactKlineData 已按列组织 (字段→数组),
直接用列数组建 polars 帧 — 无 pandas 中转、无逐行转换、无 _resolve_names
名称解析, 全市场单轮本地转换 ~1.1s → ~0.1s。时区/dtype/列序契约
复用 _normalize_minute (timestamp→北京墙钟 datetime, 输出 canonical 8 列)。
"""
tf = get_client()
try:
raw = tf.klines.intraday_universe(universe, count=count, as_dataframe=True)
raw = tf.klines.intraday_universe(universe, count=count, as_dataframe=False)
except Exception as e:
logger.warning("intraday universe fetch failed (%s): %s", universe, e)
return (pl.DataFrame(), 0)
frames = _normalize_intraday_raw(raw)
if not frames:
seg = _compact_klines_to_df(raw)
if seg.is_empty():
return (pl.DataFrame(), 0)
return (pl.concat(frames, how="diagonal_relaxed"), 1)
return (_normalize_minute(seg), 1)
def fetch_minute_single(
@@ -1050,18 +1094,13 @@ def fetch_minute_single(
end_time=_datetime_to_ms(end_time),
count=10000,
adjust="forward",
as_dataframe=True, show_progress=False,
as_dataframe=False, show_progress=False,
)
except Exception as e:
logger.warning("fetch_minute_single(%s, %s) failed: %s", symbol, trade_date, e)
return pl.DataFrame()
if isinstance(raw, dict):
sub = raw.get(symbol)
return _normalize_minute(sub) if sub is not None and len(sub) > 0 else pl.DataFrame()
if raw is not None and len(raw) > 0:
return _normalize_minute(raw)
return pl.DataFrame()
return _normalize_minute(_compact_klines_to_df(raw, default_symbol=symbol))
def fetch_adj_factor_single(symbol: str) -> pl.DataFrame:
@@ -1072,7 +1111,7 @@ def fetch_adj_factor_single(symbol: str) -> pl.DataFrame:
"""
tf = get_client()
try:
raw = tf.klines.ex_factors([symbol], as_dataframe=True, show_progress=False)
raw = tf.klines.ex_factors([symbol], as_dataframe=False, show_progress=False)
except Exception as e: # noqa: BLE001
logger.warning("fetch_adj_factor_single(%s) failed: %s", symbol, e)
return pl.DataFrame()
+4 -2
View File
@@ -44,9 +44,11 @@ import polars as pl
from app.market_time import cn_now, cn_today, in_continuous_session
from app.services import preferences
# 轮询间隔允许范围 (秒): 稳态轮单请求无并发脉冲, 下限 3s; 上限防误配。
# 轮询间隔允许范围 (秒): 稳态轮单请求无并发脉冲, 下限 3s;
# 上限 120s — universe 端点每标的只回最新 3 根, 间隔超过 3 分钟必留缺口,
# 每轮都会触发修复轮, 稳态设计失效, 故不允许配到 120s 以上。
REFRESH_INTERVAL_MIN = 3
REFRESH_INTERVAL_MAX = 300
REFRESH_INTERVAL_MAX = 120
# 等待步长 (秒): 循环小步睡眠, 便于快速停止与偏好热生效。
_LOOP_STEP_S = 2.0
# 当日覆盖滞后超过该分钟数 (≈ universe 单请求 3 根余量) → 触发全天修复轮。
+3 -3
View File
@@ -197,7 +197,7 @@ def get_minute_sync_segment_days() -> int:
# 全天修复轮 (intraday.batch 28 块爆发) 的 rpm 安全与间隔无关, 由轮次
# 调度 max(间隔, 单轮完成) 天然防重叠。
_MINUTE_REFRESH_INTERVAL_MIN = 3
_MINUTE_REFRESH_INTERVAL_MAX = 300
_MINUTE_REFRESH_INTERVAL_MAX = 120
def get_minute_refresh_enabled() -> bool:
@@ -206,7 +206,7 @@ def get_minute_refresh_enabled() -> bool:
def get_minute_refresh_interval() -> int:
"""盘中分钟增量刷新间隔(秒)。默认 6,范围 [3, 300]。"""
"""盘中分钟增量刷新间隔(秒)。默认 6,范围 [3, 120]。"""
return max(
_MINUTE_REFRESH_INTERVAL_MIN,
min(_MINUTE_REFRESH_INTERVAL_MAX, int(load().get("minute_refresh_interval", 6))),
@@ -926,7 +926,7 @@ def set_realtime_monitor_config(cfg: dict) -> dict:
if "minute_refresh_enabled" in cfg:
updates["minute_refresh_enabled"] = bool(cfg["minute_refresh_enabled"])
if "minute_refresh_interval" in cfg:
# clamp 到 [3, 300], 与 getter 一致, 防前端传越界值
# clamp 到 [3, 120], 与 getter 一致, 防前端传越界值
updates["minute_refresh_interval"] = max(
_MINUTE_REFRESH_INTERVAL_MIN,
min(_MINUTE_REFRESH_INTERVAL_MAX, int(cfg["minute_refresh_interval"])))
@@ -16,17 +16,18 @@ def _capset(batch: int = 2) -> CapabilitySet:
return CapabilitySet({Cap.INTRADAY_BATCH: CapabilityLimits(rpm=60, batch=batch)})
def _frame() -> pl.DataFrame:
def _frame() -> dict:
# _normalize_minute 的最小输入: 毫秒 timestamp → 北京墙钟 datetime。
# 时间必须落在交易时段 (时区契约守卫会拒绝非交易小时的脏数据)
from datetime import datetime
from zoneinfo import ZoneInfo
ts = int(datetime(2026, 8, 28, 9, 31, tzinfo=ZoneInfo("Asia/Shanghai")).timestamp() * 1000)
return pl.DataFrame({
# as_dataframe=False 的最小 CompactKlineData (字段 → 列数组)
return {
"timestamp": [ts],
"open": [1.0], "high": [1.0], "low": [1.0], "close": [1.0],
"volume": [100.0], "amount": [100.0],
})
}
class _FakeKlines:
@@ -167,15 +167,18 @@ def test_intraday_batch_provider_is_normalized_without_network(monkeypatch):
def intraday_batch(self, symbols, count, as_dataframe, show_progress, batch_size):
assert symbols == ["600000.SH"]
assert count == 300
assert as_dataframe is True
assert as_dataframe is False
assert show_progress is False
assert batch_size == 20
return pl.DataFrame({
"symbol": symbols,
"datetime": [datetime(2026, 7, 17, 9, 30)],
# as_dataframe=False: dict[symbol → CompactKlineData (字段→列数组)]
ts = int(datetime(2026, 7, 17, 9, 30, tzinfo=CN_TZ).timestamp() * 1000)
return {
"600000.SH": {
"timestamp": [ts],
"open": [10.0], "high": [10.1], "low": [9.9], "close": [10.0],
"volume": [1.0], "amount": [1000.0],
})
},
}
class FakeClient:
klines = FakeKlines()
+3 -3
View File
@@ -272,12 +272,12 @@ def test_refresh_preferences_defaults_and_clamp(tmp_path, monkeypatch):
preferences.save({"minute_refresh_interval": 1})
assert preferences.get_minute_refresh_interval() == 3 # 下限
preferences.save({"minute_refresh_interval": 999})
assert preferences.get_minute_refresh_interval() == 300 # 上限
assert preferences.get_minute_refresh_interval() == 120 # 上限
preferences.save({"minute_refresh_interval": 15})
assert preferences.get_minute_refresh_interval() == 15
def test_realtime_monitor_config_owns_refresh_keys(tmp_path, monkeypatch):
"""盘中增量配置归属实时监控端点 (set_realtime_monitor_config), 并 clamp 到 [3,300]。"""
"""盘中增量配置归属实时监控端点 (set_realtime_monitor_config), 并 clamp 到 [3,120]。"""
_isolated_prefs(tmp_path, monkeypatch)
saved = preferences.set_realtime_monitor_config({
"minute_refresh_enabled": True,
@@ -286,7 +286,7 @@ def test_realtime_monitor_config_owns_refresh_keys(tmp_path, monkeypatch):
assert saved["minute_refresh_enabled"] is True
assert saved["minute_refresh_interval"] == 3
saved = preferences.set_realtime_monitor_config({"minute_refresh_interval": 400})
assert saved["minute_refresh_interval"] == 300
assert saved["minute_refresh_interval"] == 120
saved = preferences.set_realtime_monitor_config({"minute_refresh_interval": 6})
assert saved["minute_refresh_interval"] == 6
+3 -3
View File
@@ -56,7 +56,7 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
const intradayInterval = prefs?.minute_intraday_refresh_interval ?? 6
// 滑块本地草稿: 拖动时即时反馈, 停顿 2s 后落库 (与行情轮询滑块一致)
const [intradayIntervalDraft, setIntradayIntervalDraft] = useState(intradayInterval)
// 盘中分钟增量 (Expert 专有): 间隔 (秒), 与后端 [3,300] clamp 对齐; 默认 6
// 盘中分钟增量 (Expert 专有): 间隔 (秒), 与后端 [3,120] clamp 对齐; 默认 6
const minuteRefreshInterval = prefs?.minute_refresh_interval ?? 6
const [minuteRefreshIntervalDraft, setMinuteRefreshIntervalDraft] = useState(minuteRefreshInterval)
// 盘中增量服务状态 (15s 轮询; 无服务时 available=false)
@@ -487,7 +487,7 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
<input
type="range"
min={3}
max={300}
max={120}
step={3}
value={minuteRefreshIntervalDraft}
disabled={!hasFullMinuteCap}
@@ -495,7 +495,7 @@ export function SettingsMonitoringPanel({ highlight }: { highlight?: string } =
className="flex-1 h-1 accent-accent cursor-pointer disabled:opacity-40 disabled:cursor-not-allowed"
/>
<span className="text-[10px] text-muted shrink-0">
{minuteRefreshIntervalDraft !== minuteRefreshInterval ? '2秒后保存' : '3s — 300s'}
{minuteRefreshIntervalDraft !== minuteRefreshInterval ? '2秒后保存' : '3s — 120s'}
</span>
</div>
{rs?.available && rs.rounds != null && rs.rounds > 0 && (