mirror of
https://ghfast.top/https://github.com/aeroxw/easy-tdx.git
synced 2026-09-12 18:04:16 +08:00
release: v1.17.2 — QFQ 深层历史负价修复
通达信服务端 QFQ 模式对长期重度除权股票(如 601088)深层历史页 返回负价格,导致回测总收益 -3087%、回撤 326.85%、年化 nan、 bollinger 崩溃、10 策略 invalid-value-in-scalar-power、MyTT divide-by-zero。 客户端兜底:检测 QFQ 负价时用 NONE+XDXR 本地重算前复权 (因子以除权日前一交易日含权收盘价为基准,保证除权日前后连续)。 同步+异步双路径一致修复,失败降级返回原值。 - 新增 src/easy_tdx/mac/adjust.py(纯函数 compute_forward_factor/ apply_forward_adjust/has_bad_prices) - MacClient/AsyncMacClient 触发本地重算,XDXR 按 (market,code) 缓存 - tests: +20 例(16 纯函数 + 4 集成),844 全绿,ruff/mypy 通过
This commit is contained in:
@@ -0,0 +1,174 @@
|
||||
"""本地前复权(QFQ)重算。
|
||||
|
||||
通达信 MAC 服务端在 QFQ 模式下,对长期重度除权股票的深层历史页会返回
|
||||
负价格(上游缺陷)。本模块用 NONE(未复权)K 线 + XDXR(除权除息)记录
|
||||
在客户端本地重算前复权序列,作为服务端 QFQ 异常时的兜底。
|
||||
|
||||
公式(见 ``examples/06_finance/xdxr_info.py``)::
|
||||
|
||||
复权价 = (原价 - 每股分红 + 每股配股价 × 每股配股比例) /
|
||||
(1 + 每股送转股比例 + 每股配股比例)
|
||||
|
||||
约定:以除权日**前一交易日**的收盘价(含权价 ``P_cum``)作为基准,前复权
|
||||
因子把该日及之前的价格乘以::
|
||||
|
||||
f = (P_cum - fenhong + peigujia × peigu) / (P_cum × (1 + songzhuangu + peigu))
|
||||
|
||||
这样调整后的价格在除权日前后连续(除权日开盘价 ≈ 含权收盘价 - 分红)。
|
||||
最新价不动(锚定最新)。fenhong/songzhuangu/peigu 为每股单位,peigujia 为元/股。
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
import numpy as np
|
||||
import pandas as pd
|
||||
|
||||
_logger = logging.getLogger(__name__)
|
||||
|
||||
# 前复权会同比缩放的 OHLC 列名(vol/amount 不动)
|
||||
_OHLC_COLS = ("open", "high", "low", "close")
|
||||
|
||||
|
||||
def compute_forward_factor(
|
||||
cum_close: float,
|
||||
fenhong: float,
|
||||
peigujia: float,
|
||||
songzhuangu: float,
|
||||
peigu: float,
|
||||
) -> float:
|
||||
"""计算单次除权除息事件的前复权乘子。
|
||||
|
||||
Args:
|
||||
cum_close: 除权日前一交易日的 NONE 收盘价(含权价)。
|
||||
fenhong: 每股分红(元)。
|
||||
peigujia: 每股配股价(元/股)。
|
||||
songzhuangu: 每股送转股比例(如 0.1 = 10 送/转 1)。
|
||||
peigu: 每股配股比例。
|
||||
|
||||
Returns:
|
||||
前复权因子。若输入非法(cum_close<=0、分母为 0、结果非有限)返回 NaN。
|
||||
"""
|
||||
if cum_close <= 0:
|
||||
return float("nan")
|
||||
denom = cum_close * (1.0 + songzhuangu + peigu)
|
||||
if denom == 0:
|
||||
return float("nan")
|
||||
factor = (cum_close - fenhong + peigujia * peigu) / denom
|
||||
if not np.isfinite(factor):
|
||||
return float("nan")
|
||||
return float(factor)
|
||||
|
||||
|
||||
def apply_forward_adjust(
|
||||
df: pd.DataFrame,
|
||||
xdxr_df: pd.DataFrame,
|
||||
) -> pd.DataFrame:
|
||||
"""对 NONE K 线应用前复权,返回新的 DataFrame。
|
||||
|
||||
遍历 XDXR 中 category==1(除权除息)的事件,按日期升序,把每个事件
|
||||
的因子累乘到该除权日**前一交易日及之前**所有 bar 的 OHLC。最新价锚定不动。
|
||||
|
||||
Args:
|
||||
df: NONE K 线,必须含 ``datetime`` 列与 OHLC 列。
|
||||
xdxr_df: ``get_xdxr_info`` 返回的 DataFrame,含 ``date``、``category``、
|
||||
``fenhong``、``peigujia``、``songzhuangu``、``peigu`` 列。
|
||||
|
||||
Returns:
|
||||
前复权后的 DataFrame(与输入同形状、同列、同索引)。无事件或异常时
|
||||
原样返回。
|
||||
"""
|
||||
out = df.copy()
|
||||
ohlc_cols = [c for c in _OHLC_COLS if c in out.columns]
|
||||
if not ohlc_cols or "datetime" not in out.columns or xdxr_df is None or xdxr_df.empty:
|
||||
return out
|
||||
|
||||
# 统一 datetime 为 pandas Timestamp(升序)
|
||||
dt = pd.to_datetime(out["datetime"])
|
||||
if not dt.is_monotonic_increasing:
|
||||
order = np.argsort(dt.to_numpy())
|
||||
out = out.iloc[order].reset_index(drop=True)
|
||||
dt = pd.to_datetime(out["datetime"])
|
||||
dt_arr = dt.to_numpy()
|
||||
|
||||
# 筛选 category==1 且至少有一个非空除权字段的事件
|
||||
if "category" not in xdxr_df.columns or "date" not in xdxr_df.columns:
|
||||
return out
|
||||
cat1 = xdxr_df[xdxr_df["category"] == 1]
|
||||
events: list[tuple[pd.Timestamp, float, float, float, float]] = []
|
||||
for _, r in cat1.iterrows():
|
||||
fh = _to_float(r.get("fenhong"))
|
||||
pjk = _to_float(r.get("peigujia"))
|
||||
sz = _to_float(r.get("songzhuangu"))
|
||||
pg = _to_float(r.get("peigu"))
|
||||
if fh is None and pjk is None and sz is None and pg is None:
|
||||
continue
|
||||
try:
|
||||
ed = pd.Timestamp(str(r["date"]))
|
||||
except (ValueError, TypeError):
|
||||
continue
|
||||
events.append((ed, fh or 0.0, pjk or 0.0, sz or 0.0, pg or 0.0))
|
||||
if not events:
|
||||
return out
|
||||
events.sort(key=lambda e: e[0])
|
||||
|
||||
# 取除权日前一交易日的 NONE 收盘价(cum-div close)作为基准
|
||||
none_close = out["close"].to_numpy(dtype=float) if "close" in out.columns else None
|
||||
|
||||
for col in ohlc_cols:
|
||||
arr = out[col].to_numpy(dtype=float).copy()
|
||||
for ed, fh, pjk, sz, pg in events:
|
||||
# searchsorted(>=): 第一个 >= ex-date 的位置;其前一根即为含权收盘
|
||||
idx = int(np.searchsorted(dt_arr, np.datetime64(ed), side="left"))
|
||||
cum_idx = idx - 1
|
||||
if cum_idx < 0 or cum_idx >= len(arr):
|
||||
continue
|
||||
if none_close is not None:
|
||||
cum_close = float(none_close[cum_idx])
|
||||
else:
|
||||
cum_close = float(arr[cum_idx])
|
||||
factor = compute_forward_factor(cum_close, fh, pjk, sz, pg)
|
||||
if not np.isfinite(factor):
|
||||
_logger.warning(
|
||||
"QFQ 本地重算:事件 %s 因子非法(cum_close=%s fh=%s sz=%s pg=%s),跳过",
|
||||
ed.date(), cum_close, fh, sz, pg,
|
||||
)
|
||||
continue
|
||||
arr[: cum_idx + 1] *= factor
|
||||
out[col] = arr
|
||||
|
||||
return out
|
||||
|
||||
|
||||
def _to_float(v: object) -> float | None:
|
||||
"""安全转 float,None/NaN 返回 None。"""
|
||||
if v is None:
|
||||
return None
|
||||
try:
|
||||
f = float(v) # type: ignore[arg-type]
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
if not np.isfinite(f):
|
||||
return None
|
||||
return f
|
||||
|
||||
|
||||
def has_bad_prices(df: pd.DataFrame) -> bool:
|
||||
"""检测 QFQ 结果是否含非法价格(<=0 或非有限值)。
|
||||
|
||||
Args:
|
||||
df: 待检测的 K 线 DataFrame。
|
||||
|
||||
Returns:
|
||||
任一 OHLC 列含 <=0 或 NaN/inf 时返回 True。
|
||||
"""
|
||||
for col in _OHLC_COLS:
|
||||
if col not in df.columns:
|
||||
continue
|
||||
arr = df[col].to_numpy(dtype=float)
|
||||
if not np.all(np.isfinite(arr)):
|
||||
return True
|
||||
if np.any(arr <= 0):
|
||||
return True
|
||||
return False
|
||||
+191
-25
@@ -153,6 +153,9 @@ class MacClient:
|
||||
self._auto_reconnect = auto_reconnect
|
||||
self._heartbeat_interval = heartbeat_interval
|
||||
self._conn = TdxConnection(self._host, self._port, self._timeout)
|
||||
# XDXR(除权除息)记录缓存:(market, code) -> DataFrame。
|
||||
# 仅在服务端 QFQ 返回异常(负价)时用于本地前复权重算。
|
||||
self._xdxr_cache: dict[tuple[int, str], pd.DataFrame] = {}
|
||||
|
||||
# ------------------------------------------------------------------ #
|
||||
# 工厂方法
|
||||
@@ -327,6 +330,102 @@ class MacClient:
|
||||
|
||||
return _quotes_to_df(all_quotes)
|
||||
|
||||
# ------------------------------------------------------------------ #
|
||||
# QFQ 本地重算(服务端 QFQ 对深层历史返回负价时的兜底)
|
||||
# ------------------------------------------------------------------ #
|
||||
|
||||
def _fetch_kline_pages(
|
||||
self,
|
||||
market: int,
|
||||
code: str,
|
||||
period: Period,
|
||||
start: int,
|
||||
count: int,
|
||||
times: int,
|
||||
fq: Adjust,
|
||||
) -> list[MacBar]:
|
||||
"""分页拉取指定复权类型的 K 线(返回 oldest→newest 的 MacBar 列表)。"""
|
||||
all_bars: list[MacBar] = []
|
||||
fetched = 0
|
||||
offset = start
|
||||
while fetched < count:
|
||||
page_size = min(count - fetched, _KLINE_PAGE_SIZE)
|
||||
bars = self._execute(
|
||||
SymbolBarCmd(
|
||||
market=market,
|
||||
code=code,
|
||||
period=period,
|
||||
times=times,
|
||||
start=offset,
|
||||
count=page_size,
|
||||
fq=fq,
|
||||
)
|
||||
)
|
||||
if not bars:
|
||||
break
|
||||
all_bars = bars + all_bars
|
||||
fetched += len(bars)
|
||||
offset += len(bars)
|
||||
if len(bars) < page_size:
|
||||
break
|
||||
return all_bars
|
||||
|
||||
def _fetch_xdxr_records(self, market: int, code: str) -> pd.DataFrame | None:
|
||||
"""通过主协议客户端(TdxClient)拉取除权除息记录。
|
||||
|
||||
MAC 主机池不响应 XDXR(0x0c1f),需连 get_known_hosts 主机池。
|
||||
结果按 (market, code) 缓存。失败返回 None(调用方降级)。
|
||||
"""
|
||||
key = (market, code)
|
||||
if key in self._xdxr_cache:
|
||||
return self._xdxr_cache[key]
|
||||
try:
|
||||
# 函数内 import 避免循环依赖(client 依赖 mac,mac 不应依赖 client)
|
||||
from .. import Market
|
||||
from ..client import TdxClient
|
||||
|
||||
with TdxClient.from_best_host(timeout=self._timeout) as tc:
|
||||
xd = tc.get_xdxr_info(Market(market), code)
|
||||
except Exception as exc: # noqa: BLE001 - 降级,不中断 kline 获取
|
||||
_logger.warning(
|
||||
"QFQ 本地重算:获取 %s %s XDXR 失败,降级返回服务端 QFQ:%s", market, code, exc,
|
||||
)
|
||||
return None
|
||||
if xd is None or xd.empty:
|
||||
return None
|
||||
self._xdxr_cache[key] = xd
|
||||
return xd
|
||||
|
||||
def _local_recompute_qfq(
|
||||
self,
|
||||
df: pd.DataFrame,
|
||||
market: int,
|
||||
code: str,
|
||||
) -> pd.DataFrame:
|
||||
"""对 QFQ 异常的 K 线用 NONE + XDXR 本地重算前复权。
|
||||
|
||||
Args:
|
||||
df: 服务端 QFQ 结果(含异常)。
|
||||
market: 市场代码。
|
||||
code: 股票代码。
|
||||
|
||||
Returns:
|
||||
重算后的 DataFrame;XDXR 取不到或重算仍异常时原样返回 df。
|
||||
"""
|
||||
from .adjust import apply_forward_adjust, has_bad_prices
|
||||
|
||||
xd = self._fetch_xdxr_records(market, code)
|
||||
if xd is None:
|
||||
return df
|
||||
out = apply_forward_adjust(df, xd)
|
||||
if has_bad_prices(out):
|
||||
_logger.warning("QFQ 本地重算后 %s %s 仍含非法价格,降级返回服务端 QFQ", market, code)
|
||||
return df
|
||||
_logger.warning(
|
||||
"QFQ 本地重算:%s %s 服务端深层历史返回负价,已用 NONE+XDXR 重算前复权", market, code,
|
||||
)
|
||||
return out
|
||||
|
||||
# ------------------------------------------------------------------ #
|
||||
# K 线(支持复权)
|
||||
# ------------------------------------------------------------------ #
|
||||
@@ -358,32 +457,21 @@ class MacClient:
|
||||
(= 开始 + 周期时长,与 Tushare/同花顺对齐,上午最后一根标 11:30)。
|
||||
仅对分钟级周期生效;日线及以上不受影响。
|
||||
"""
|
||||
all_bars: list[MacBar] = []
|
||||
fetched = 0
|
||||
offset = start
|
||||
|
||||
while fetched < count:
|
||||
page_size = min(count - fetched, _KLINE_PAGE_SIZE)
|
||||
bars = self._execute(
|
||||
SymbolBarCmd(
|
||||
market=market,
|
||||
code=code,
|
||||
period=period,
|
||||
times=times,
|
||||
start=offset,
|
||||
count=page_size,
|
||||
fq=adjust,
|
||||
)
|
||||
)
|
||||
if not bars:
|
||||
break
|
||||
all_bars = bars + all_bars
|
||||
fetched += len(bars)
|
||||
offset += len(bars)
|
||||
if len(bars) < page_size:
|
||||
break
|
||||
|
||||
all_bars = self._fetch_kline_pages(market, code, period, start, count, times, adjust)
|
||||
df = _to_df(all_bars)
|
||||
|
||||
# QFQ 兜底:服务端对深层历史可能返回负价/零价,此时用 NONE+XDXR 本地重算。
|
||||
if adjust == Adjust.QFQ and not df.empty:
|
||||
from .adjust import has_bad_prices
|
||||
|
||||
if has_bad_prices(df):
|
||||
none_bars = self._fetch_kline_pages(
|
||||
market, code, period, start, count, times, Adjust.NONE
|
||||
)
|
||||
df = _to_df(none_bars) if none_bars else df
|
||||
if not df.empty:
|
||||
df = self._local_recompute_qfq(df, market, code)
|
||||
|
||||
delta = _period_to_minutes(period, times)
|
||||
is_intraday = delta is not None
|
||||
return _apply_bar_time_align_df(
|
||||
@@ -1071,6 +1159,9 @@ class AsyncMacClient(AsyncHeartbeatMixin):
|
||||
self._conn = AsyncTdxConnection(self._host, self._port, self._timeout)
|
||||
self._execute_lock = asyncio.Lock()
|
||||
self._heartbeat_task: asyncio.Task[None] | None = None
|
||||
# XDXR(除权除息)记录缓存:(market, code) -> DataFrame。
|
||||
# 仅在服务端 QFQ 返回异常(负价)时用于本地前复权重算。
|
||||
self._xdxr_cache: dict[tuple[int, str], pd.DataFrame] = {}
|
||||
|
||||
# ------------------------------------------------------------------ #
|
||||
# 工厂方法
|
||||
@@ -1233,6 +1324,50 @@ class AsyncMacClient(AsyncHeartbeatMixin):
|
||||
|
||||
return _quotes_to_df(all_quotes)
|
||||
|
||||
# ------------------------------------------------------------------ #
|
||||
# QFQ 本地重算(服务端 QFQ 对深层历史返回负价时的兜底)
|
||||
# 同步方法:XDXR 经 TdxClient(同步主协议)获取,由 asyncio.to_thread 调用。
|
||||
# ------------------------------------------------------------------ #
|
||||
|
||||
def _fetch_xdxr_records(self, market: int, code: str) -> pd.DataFrame | None:
|
||||
"""通过主协议客户端(TdxClient)拉取除权除息记录(同 MacClient)。"""
|
||||
key = (market, code)
|
||||
if key in self._xdxr_cache:
|
||||
return self._xdxr_cache[key]
|
||||
try:
|
||||
from .. import Market
|
||||
from ..client import TdxClient
|
||||
|
||||
with TdxClient.from_best_host(timeout=self._timeout) as tc:
|
||||
xd = tc.get_xdxr_info(Market(market), code)
|
||||
except Exception as exc: # noqa: BLE001 - 降级,不中断 kline 获取
|
||||
_logger.warning(
|
||||
"QFQ 本地重算:获取 %s %s XDXR 失败,降级返回服务端 QFQ:%s", market, code, exc,
|
||||
)
|
||||
return None
|
||||
if xd is None or xd.empty:
|
||||
return None
|
||||
self._xdxr_cache[key] = xd
|
||||
return xd
|
||||
|
||||
def _local_recompute_qfq(
|
||||
self, df: pd.DataFrame, market: int, code: str,
|
||||
) -> pd.DataFrame:
|
||||
"""对 QFQ 异常的 K 线用 NONE+XDXR 本地重算前复权(同 MacClient)。"""
|
||||
from .adjust import apply_forward_adjust, has_bad_prices
|
||||
|
||||
xd = self._fetch_xdxr_records(market, code)
|
||||
if xd is None:
|
||||
return df
|
||||
out = apply_forward_adjust(df, xd)
|
||||
if has_bad_prices(out):
|
||||
_logger.warning("QFQ 本地重算后 %s %s 仍含非法价格,降级返回服务端 QFQ", market, code)
|
||||
return df
|
||||
_logger.warning(
|
||||
"QFQ 本地重算:%s %s 服务端深层历史返回负价,已用 NONE+XDXR 重算前复权", market, code,
|
||||
)
|
||||
return out
|
||||
|
||||
# ------------------------------------------------------------------ #
|
||||
# K 线
|
||||
# ------------------------------------------------------------------ #
|
||||
@@ -1276,6 +1411,37 @@ class AsyncMacClient(AsyncHeartbeatMixin):
|
||||
break
|
||||
|
||||
df = _to_df(all_bars)
|
||||
|
||||
# QFQ 兜底:服务端对深层历史可能返回负价/零价,此时用 NONE+XDXR 本地重算。
|
||||
if adjust == Adjust.QFQ and not df.empty:
|
||||
from .adjust import has_bad_prices
|
||||
|
||||
if has_bad_prices(df):
|
||||
# 异步重抓 NONE
|
||||
none_bars: list[MacBar] = []
|
||||
nfetched = 0
|
||||
noffset = start
|
||||
while nfetched < count:
|
||||
nps = min(count - nfetched, _KLINE_PAGE_SIZE)
|
||||
nb = await self._execute(
|
||||
SymbolBarCmd(
|
||||
market=market, code=code, period=period, times=times,
|
||||
start=noffset, count=nps, fq=Adjust.NONE,
|
||||
)
|
||||
)
|
||||
if not nb:
|
||||
break
|
||||
none_bars = nb + none_bars
|
||||
nfetched += len(nb)
|
||||
noffset += len(nb)
|
||||
if len(nb) < nps:
|
||||
break
|
||||
if none_bars:
|
||||
# XDXR 获取涉及同步网络 IO,放线程执行
|
||||
df = await asyncio.to_thread(
|
||||
self._local_recompute_qfq, _to_df(none_bars), market, code,
|
||||
)
|
||||
|
||||
delta = _period_to_minutes(period, times)
|
||||
is_intraday = delta is not None
|
||||
return _apply_bar_time_align_df(
|
||||
|
||||
Reference in New Issue
Block a user