feat(factor): P1-B批adapter四新域读取层+CCC三表自算 [nas]

- E08 ccc: 365×(AR均值/REV_TTM + INV均值/COGS_TTM − AP均值/COGS_TTM),
  均值=t与t−4期末平均(_mean4),ACCOUNTS_PAYABLE 入 balance 读取清单,
  金融股置 NaN(NAS 实测列名核实)
- E11 send_total_12m: dividend 回退链 PIT(Plan→PreNotice→AgmPum→Plan,
  全空行丢弃),cum(d)−cum(d−365) 两端 asof 实现开区间滚动窗;域内无事件
  →NaN 不填 0(与缺文件区分)
- E12 gdhs_chg: 自建事件表(代码,截止日)→户数,同键取最新公告,相邻事件
  Δln,公告日期 asof;不用「增减比例」列;实测截止日列名无第二连字符(双名兼容)
- E13 topholder_chg: 只读 factor_cache/top_holders_agg.parquet,期→法定
  披露截止(Q1→04-30/H1→08-31/Q3→10-31/年报→次年04-30)保守 PIT
- D14 ep_vb/bp_vb: vb 按年文件+2026-01-01~08-12 缺口 static/valuation(ak
  中文列)倒数补;实测 exchange=SH/SZ(SSE/SZSE 全称变体兼容);1/pe 守卫 0→NaN
- 四域引用列驱动加载(未引用不读盘)+缺域一次性 warning 优雅降级全 NaN
- conftest 四域合成 fixture;NAS 真数据冒烟通过(300750 ccc=12.5d/
  000333 ccc=−9.2d/600519 ep=1/32PE)
This commit is contained in:
2026-09-08 18:54:04 +08:00
parent 6e719de0c6
commit 60dd84ad4a
3 changed files with 791 additions and 5 deletions
+392 -4
View File
@@ -16,6 +16,22 @@
- 日期列为 "YYYY-MM-DD 00:00:00" 字符串;东财金额单位=元
- valuation(中文列)不读: 估值类因子市值 = close × SHARE_CAPITAL 自算(表达式层)
P1-B 四新域(NAS 实测 2026-09-08):
- dividend 按股: static/dividend/{code}.{SH|SZ}_dividend.parquet;code 列为
baostock 小写格式(sh.600519)与文件名不同 → 以文件名路由;日期/数值列为
字符串('' 为空);一行=一次分红事件,同年可多行
- gdhs 按期: static/gdhs/{YYYYMMDD}_gdhs.parquet(53 期 2013Q1→2026Q1);
用 代码/股东户数-本次/股东户数统计截止日-本次/公告日期 四列(实测列名
无第二连字符;噪声行情快照列不读)
- top_holders: 只读聚合产物 factor_cache/top_holders_agg.parquet(static
兄弟目录;由 scripts/factor_research/preaggregate_top_holders.py 一次性
预聚合 108,610 个按股×期小文件)
- valuation_baostock 按年: data/valuation_baostock/{year}.parquet(static
兄弟目录);date 为字符串,exchange 实测为 SH/SZ(映射回 SSE/SZSE);
2026 文件仅 08-13 起(bs 日喂起点)→ 01-01~08-12
缺口用 static/valuation(ak em 中文列 PE(TTM)/市净率)倒数补
- 四域缺任一 → 相关特征列全 NaN + 一次性 warning(优雅降级,本地可跑通)
输出: vt_symbol × datetime(日频) × FEATURE_COLUMNS,join 到 alpha_df 后
供 cs_rank(列) 表达式直接消费。
"""
@@ -23,7 +39,7 @@ from __future__ import annotations
import os
import warnings
from datetime import datetime
from datetime import date, datetime
import polars as pl
@@ -75,7 +91,9 @@ _BALANCE_COLS = ["TOTAL_ASSETS", "TOTAL_PARENT_EQUITY", "ACCOUNTS_RECE",
"SHORT_FIN_PAYABLE", "NONCURRENT_LIAB_1YEAR", "LONG_LOAN",
"BOND_PAYABLE", "LEASE_LIAB",
# P1 新增存量列
"INVENTORY", "MONETARYFUNDS"]
"INVENTORY", "MONETARYFUNDS",
# P1-B 新增存量列(E08 CCC 分子: 应付账款,NAS 实测列名核实)
"ACCOUNTS_PAYABLE"]
# IBD 口径钉死含 LEASE_LIAB 版(survey §3.5 E 族注意点「定稿钉死」;2019 前该列
# 整体缺失按 0,与一年内到期非流动负债等组件同款处理)
_IBD_PARTS = ["SHORT_LOAN", "SHORT_FIN_PAYABLE", "NONCURRENT_LIAB_1YEAR",
@@ -117,6 +135,12 @@ FEATURE_COLUMNS: list[str] = [
"debt_issue_ttm", "interest_cover", "goodwill_growth",
# ---- P1 批新增(F 预期 3 + SUE 严窗变体;forecast_beat 走独立兑现差事件流)----
"sue_eps", "disclosure_speed", "sue_np_strict", "forecast_beat",
# ---- P1-B 批新增(四新域 5 因子;契约由 test_fundamental_p1b_* 锁定)----
"ccc", # E08 三表自算现金转换周期
"send_total_12m", # E11 dividend 事件流(365 天滚动合计)
"gdhs_chg", # E12 股东户数相邻事件 Δln
"topholder_chg", # E13 十大流通占比相邻期 Δ(聚合产物)
"ep_vb", "bp_vb", # D14 长史 EP/BP(vb 日频域)
]
# 金融股置 NaN 的特征(盈利质量 B 族 + 成长 C 族 + 费用类,§7 红线 5;
# survey A 族注意点: A14/A16 金融无三费结构 → rd/sale_expense_ratio 同剔)
@@ -129,12 +153,16 @@ _FIN_NULL_COLS = ["tacc", "nonrec_ratio", "impairment_ratio",
"inventory_anomaly", "cash_ibd_product", "vsig", "vsig_acc",
"vsig_cfo", "da_intensity", "gm_nm_scissors", "profit_streak",
"np_accel", "rev_accel", "nm_delta", "rev_cagr5", "np_cagr5",
"nwc_growth", "invest_growth", "equity_growth"]
"nwc_growth", "invest_growth", "equity_growth",
# P1-B: E08 CCC 属营运效率,金融股无 CCC 概念(§7 红线扩展)
"ccc"]
# forecast 事件流日频列(公告日 asof): P0 两列 + P1 年化预告净利(D13)
_FORECAST_COLS = ("forecast_type_score", "forecast_change_pct",
"forecast_np_annualized")
# 预告兑现差独立事件流(F07 前视红线: 锚 = max(实际披露日, 预告公告日))
_BEAT_COLS = ("forecast_beat",)
# P1-B 四新域列(非报告期特征,由各域事件/日频流产出;报表加载须排除)
_DOMAIN_COLS = ("send_total_12m", "gdhs_chg", "topholder_chg", "ep_vb", "bp_vb")
# ==================== 读取层 ====================
@@ -395,6 +423,10 @@ def _compute_report_features(reports: pl.DataFrame) -> pl.DataFrame:
_safe_ratio(pl.col("q_np"), ta).alias("_np_ta"),
_safe_ratio(pl.col("q_np") - pl.col("q_cfo"), ta).alias("_accq_ta"),
_safe_ratio(pl.col("q_cfo"), ta).alias("_cfo_ta"),
# P1-B: E08 CCC 三分子(「均值」= t 与 t−4 报告期期末余额平均)
_mean4("ACCOUNTS_RECE").alias("_ar_avg"),
_mean4("INVENTORY").alias("_inv_avg"),
_mean4("ACCOUNTS_PAYABLE").alias("_ap_avg"),
)
# 第三级: 行本地比率(无时序,无需 over)
@@ -440,6 +472,12 @@ def _compute_report_features(reports: pl.DataFrame) -> pl.DataFrame:
pl.when(pl.col("ttm__int0") > 0)
.then(pl.col("_ebit_ttm") / pl.col("ttm__int0"))
.otherwise(None).alias("interest_cover"),
# P1-B E08 现金转换周期(三表自算,不用 abstract 域):
# DSO=365×AR均值/REV_TTM;DIO=365×INV均值/COGS_TTM;
# DPO=365×AP均值/COGS_TTM;CCC=DSO+DIODPO(缺任一原料 NaN 传播)
(_safe_ratio(pl.col("_ar_avg") * 365.0, pl.col("ttm_rev"))
+ _safe_ratio(pl.col("_inv_avg") * 365.0, pl.col("ttm_cogs"))
- _safe_ratio(pl.col("_ap_avg") * 365.0, pl.col("ttm_cogs"))).alias("ccc"),
)
# 第四级: 跨期差分/同比/剪刀差(over 组内时序)
@@ -650,6 +688,319 @@ def _load_forecast_events(codes: list[str], data_dir: str) -> tuple[pl.DataFrame
return fc_events, fc_pair
# ==================== P1-B 四新域事件层 ====================
# 缺域告警去重(每 (域, 路径) 一次;优雅降级 = 特征列全 NaN 不崩,
# 本地开发无 NAS 数据也能跑通管线)
_missing_warned: set[tuple[str, str]] = set()
# dividend PIT 日期回退链(首列空往后回退;全空行丢弃)
_DIVIDEND_CHAIN = ["dividPlanAnnounceDate", "dividPreNoticeDate",
"dividAgmPumDate", "dividPlanDate"]
# gdhs 截止日列名(NAS 实测「股东户数统计截止日-本次」;兼容任务书连字符变体)
_GDHS_CUTOFF_CANDIDATES = ["股东户数统计截止日-本次", "股东户数-统计截止日-本次"]
# vb 2026 缺口补口窗口(bs 日喂起点 2026-08-13 → 01-01~08-12 由 static/valuation 补)
_VB_GAP = (date(2026, 1, 1), date(2026, 8, 12))
def _warn_domain_missing(key: str, path: str) -> None:
if (key, path) not in _missing_warned:
_missing_warned.add((key, path))
warnings.warn(
f"[fundamental_adapter] 数据域 {key} 缺失({path}) → 相关特征列全 NaN"
f"(本地开发无 NAS 数据可忽略;NAS 真跑前先确认路径)")
def _code6_to_vt(code: pl.Expr) -> pl.Expr:
"""6 位代码 → vt_symbol(60→SSE;北交 92/43/82/83→BJSE;其余→SZSE;
与 forecast 映射同款,BJ 不落入 SZSE)."""
is_bj = (code.str.starts_with("92") | code.str.starts_with("43")
| code.str.starts_with("82") | code.str.starts_with("83"))
return (pl.when(code.str.starts_with("60")).then(code + pl.lit(".SSE"))
.when(is_bj).then(code + pl.lit(".BJSE"))
.otherwise(code + pl.lit(".SZSE")))
def _load_dividend_cum(codes: list[str], data_dir: str) -> pl.DataFrame:
"""dividend 按股文件 → (vt_symbol, eff, _send_cum) 累计事件流.
事件 = 一次分红;eff = 回退链首个非空公告日(全空行丢弃);值 =
dividStocksPs(每股送转合计)。按股日序 cum_sum 供 E11 的 365 天滚动窗
两端口径相减(cum(d) cum(d365) = 窗 (d365, d] 内事件合计)。
code 列为 baostock 小写格式(sh.600519)与文件名不同 → 以文件名路由
(与三表读取同模式,列内容不参与 join)。
"""
schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "_send_cum": pl.Float64}
div_dir = os.path.join(data_dir, "dividend")
if not os.path.isdir(div_dir):
_warn_domain_missing("dividend", div_dir)
return pl.DataFrame(schema=schema)
frames = []
for vt in codes:
file_code = _vt_to_file_code(vt)
if file_code is None:
continue
path = os.path.join(div_dir, f"{file_code}_dividend.parquet")
if not os.path.exists(path):
continue
try:
df = pl.read_parquet(path)
except Exception:
continue
chain = [c for c in _DIVIDEND_CHAIN if c in df.columns]
if not chain or "dividStocksPs" not in df.columns:
continue
# 日期列字符串('' 为空;兼容 'YYYY-MM-DD 00:00:00')→ coalesce 回退链
eff = pl.coalesce([
pl.col(c).cast(pl.Utf8).str.slice(0, 10).str.to_date("%Y-%m-%d", strict=False)
for c in chain])
df = df.select(
pl.lit(vt).alias("vt_symbol"),
eff.alias("eff"),
pl.col("dividStocksPs").cast(pl.Float64, strict=False).alias("_send"),
).filter(pl.col("eff").is_not_null())
if df.height:
frames.append(df)
if not frames:
return pl.DataFrame(schema=schema)
ev = pl.concat(frames).sort(["vt_symbol", "eff"])
return ev.with_columns(
pl.col("_send").cum_sum().over(_SYM).alias("_send_cum")
).select(["vt_symbol", "eff", "_send_cum"]).with_columns(
pl.col("eff").cast(pl.Datetime("us")))
def _load_gdhs_events(codes: list[str], data_dir: str) -> pl.DataFrame:
"""gdhs 按期文件 → (vt_symbol, eff, gdhs_chg) 事件流.
事件表: (代码, 统计截止日)→股东户数;同键多行取最新公告(修订口径);
gdhs_chg = 相邻事件(按截止日排序)Δln(户数),首事件无上期 → NaN;
PIT = 公告日期。红线: 不用「股东户数-增减比例」列(每股截止日不规则,
非统一季环比)。
"""
schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "gdhs_chg": pl.Float64}
gd_dir = os.path.join(data_dir, "gdhs")
if not os.path.isdir(gd_dir):
_warn_domain_missing("gdhs", gd_dir)
return pl.DataFrame(schema=schema)
code_set = set(codes)
frames = []
for fname in sorted(os.listdir(gd_dir)):
if not fname.endswith(".parquet"):
continue
try:
f = pl.read_parquet(os.path.join(gd_dir, fname))
except Exception:
continue
cutoff = next((c for c in _GDHS_CUTOFF_CANDIDATES if c in f.columns), None)
if cutoff is None or not all(
c in f.columns for c in ("代码", "股东户数-本次", "公告日期")):
continue
code = pl.col("代码").cast(pl.Utf8).str.strip_chars().str.zfill(6)
f = f.select(
_code6_to_vt(code).alias("vt_symbol"),
pl.col("股东户数-本次").cast(pl.Float64, strict=False).alias("_holders"),
pl.col(cutoff).cast(pl.Date, strict=False).alias("_cutoff"),
pl.col("公告日期").cast(pl.Date, strict=False).alias("_ann"),
).filter(pl.col("vt_symbol").is_in(code_set) & pl.col("_ann").is_not_null()
& pl.col("_cutoff").is_not_null() & pl.col("_holders").is_not_null())
if f.height:
frames.append(f)
if not frames:
return pl.DataFrame(schema=schema)
ev = (pl.concat(frames)
.sort(["vt_symbol", "_cutoff", "_ann"])
.group_by(["vt_symbol", "_cutoff"]).agg( # 同 (股, 截止日) 取最新公告
pl.col("_holders").last(), pl.col("_ann").last())
.sort(["vt_symbol", "_cutoff"]))
return ev.with_columns(
pl.col("_holders").shift(1).over(_SYM).alias("_prev")
).with_columns(
pl.when((pl.col("_prev") > 0) & (pl.col("_holders") > 0))
.then((pl.col("_holders") / pl.col("_prev")).log())
.otherwise(None).alias("gdhs_chg")
).select(
pl.col("vt_symbol"),
pl.col("_ann").cast(pl.Datetime("us")).alias("eff"),
pl.col("gdhs_chg"),
).sort("eff")
def _load_topholder_events(codes: list[str], data_dir: str) -> pl.DataFrame:
"""top_holders 聚合产物 → (vt_symbol, eff, topholder_chg) 事件流.
聚合产物 = scripts/factor_research/preaggregate_top_holders.py 输出
(file_code, period, hold_pct, n_holders),置于 static 兄弟目录
factor_cache/ 下(绝不写 static 树)。期→PIT 无法定披露日 → 法定披露
截止近似(保守侧): Q1→04-30 / H1→08-31 / Q3→10-31 / 年报→次年04-30。
topholder_chg = 最近期十大合计占比 − 上期(首期无上期 → NaN)。
"""
schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "topholder_chg": pl.Float64}
agg_path = os.path.join(os.path.dirname(os.path.abspath(data_dir)),
"factor_cache", "top_holders_agg.parquet")
if not os.path.exists(agg_path):
_warn_domain_missing("top_holders_agg", agg_path)
return pl.DataFrame(schema=schema)
try:
agg = pl.read_parquet(agg_path)
except Exception:
_warn_domain_missing("top_holders_agg", agg_path)
return pl.DataFrame(schema=schema)
if not all(c in agg.columns for c in ("file_code", "period", "hold_pct")):
return pl.DataFrame(schema=schema)
parts = pl.col("file_code").str.split(".")
ev = agg.with_columns(
parts.list.get(0).alias("_c"), parts.list.get(1).alias("_x"),
pl.col("period").cast(pl.Date, strict=False),
pl.col("hold_pct").cast(pl.Float64, strict=False),
).with_columns(
pl.when(pl.col("_x") == "SH").then(pl.col("_c") + pl.lit(".SSE"))
.otherwise(pl.col("_c") + pl.lit(".SZSE")).alias("vt_symbol")
).filter(
pl.col("vt_symbol").is_in(set(codes)) & pl.col("period").is_not_null()
& pl.col("hold_pct").is_not_null()
).sort(["vt_symbol", "period"])
m = pl.col("period").dt.month()
y = pl.col("period").dt.year()
deadline = (
pl.when(m == 3).then(pl.date(y, 4, 30))
.when(m == 6).then(pl.date(y, 8, 31))
.when(m == 9).then(pl.date(y, 10, 31))
.otherwise(pl.date(y + 1, 4, 30))
)
return ev.with_columns(
pl.col("hold_pct").shift(1).over(_SYM).alias("_prev")
).with_columns(
(pl.col("hold_pct") - pl.col("_prev")).alias("topholder_chg")
).select(
pl.col("vt_symbol"),
deadline.cast(pl.Datetime("us")).alias("eff"),
pl.col("topholder_chg"),
).sort("eff")
def _inv_guard(col: str) -> pl.Expr:
"""1/x 守卫: x 缺失/为 0 → NaN(亏损 baostock pe=NaN 天然 NaN)."""
return (pl.when(pl.col(col).is_not_null() & (pl.col(col) != 0))
.then(1.0 / pl.col(col)).otherwise(None))
def _load_vb_daily(codes: list[str], data_dir: str,
start: str, end: str) -> pl.DataFrame:
"""valuation_baostock 按年文件(+static/valuation 补 2026 缺口) →
(vt_symbol, eff, ep_vb, bp_vb) 日频流.
vb 在 static 的兄弟目录(data/valuation_baostock/{year}.parquet),
date 为字符串须转 Date;ep=1/peTTM、bp=1/pbMRQ。2026-01-01~08-12
缺口(bs 日喂 08-13 起)用 static/valuation(ak em 中文列:
PE(TTM)/市净率)倒数补——两源亏损口径不同(baostock NaN/ak 可为负),
倒数后均按原始符号保留,分域验证时注意。
"""
schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"),
"ep_vb": pl.Float64, "bp_vb": pl.Float64}
root = os.path.dirname(os.path.abspath(data_dir))
vb_dir = os.path.join(root, "valuation_baostock")
if not os.path.isdir(vb_dir):
_warn_domain_missing("valuation_baostock", vb_dir)
return pl.DataFrame(schema=schema)
code_set = set(codes)
d_lo = datetime.strptime(start, "%Y-%m-%d").date()
d_hi = datetime.strptime(end, "%Y-%m-%d").date()
frames = []
for year in range(d_lo.year, d_hi.year + 1):
path = os.path.join(vb_dir, f"{year}.parquet")
if not os.path.exists(path):
continue
try:
f = pl.read_parquet(path)
except Exception:
continue
if not all(c in f.columns for c in
("symbol", "exchange", "date", "peTTM", "pbMRQ")):
continue
raw_date = pl.col("date")
d = (raw_date.cast(pl.Utf8).str.slice(0, 10).str.to_date("%Y-%m-%d", strict=False)
if f.schema["date"] == pl.Utf8 else raw_date.cast(pl.Date, strict=False))
# exchange 实测为 SH/SZ(NAS 2026-09-08,任务书 SSE/SZSE 变体兼容)
xsuf = (pl.when(pl.col("exchange") == pl.lit("SH")).then(pl.lit("SSE"))
.when(pl.col("exchange") == pl.lit("SZ")).then(pl.lit("SZSE"))
.otherwise(pl.col("exchange")))
f = f.select(
(pl.col("symbol").cast(pl.Utf8).str.strip_chars() + pl.lit(".")
+ xsuf.cast(pl.Utf8).str.strip_chars()).alias("vt_symbol"),
d.alias("eff"),
pl.col("peTTM").cast(pl.Float64, strict=False).alias("_pe"),
pl.col("pbMRQ").cast(pl.Float64, strict=False).alias("_pb"),
).filter(pl.col("vt_symbol").is_in(code_set) & pl.col("eff").is_not_null()
& (pl.col("eff") >= d_lo) & (pl.col("eff") <= d_hi))
if f.height:
frames.append(f)
# 2026 缺口补口: static/valuation 按股中文列
gap_lo, gap_hi = max(_VB_GAP[0], d_lo), min(_VB_GAP[1], d_hi)
if gap_lo <= gap_hi:
val_dir = os.path.join(data_dir, "valuation")
if not os.path.isdir(val_dir):
_warn_domain_missing("static/valuation(2026 补口)", val_dir)
else:
for vt in codes:
file_code = _vt_to_file_code(vt)
if file_code is None:
continue
path = os.path.join(val_dir, f"{file_code}_valuation.parquet")
if not os.path.exists(path):
continue
try:
f = pl.read_parquet(path)
except Exception:
continue
have = [c for c in ("数据日期", "PE(TTM)", "市净率") if c in f.columns]
if "数据日期" not in have:
continue
f = _norm_dates(f.select(have), ["数据日期"])
pe = (pl.col("PE(TTM)").cast(pl.Float64, strict=False)
if "PE(TTM)" in have else pl.lit(None, pl.Float64))
pb = (pl.col("市净率").cast(pl.Float64, strict=False)
if "市净率" in have else pl.lit(None, pl.Float64))
f = f.select(
pl.lit(vt).alias("vt_symbol"),
pl.col("数据日期").alias("eff"),
pe.alias("_pe"), pb.alias("_pb"),
).filter(pl.col("eff").is_not_null()
& (pl.col("eff") >= gap_lo) & (pl.col("eff") <= gap_hi))
if f.height:
frames.append(f)
if not frames:
return pl.DataFrame(schema=schema)
out = (pl.concat(frames).sort(["vt_symbol", "eff"])
.unique(subset=["vt_symbol", "eff"], keep="last")
.with_columns(
_inv_guard("_pe").alias("ep_vb"),
_inv_guard("_pb").alias("bp_vb"),
pl.col("eff").cast(pl.Datetime("us"))))
return out.select(["vt_symbol", "eff", "ep_vb", "bp_vb"]).sort("eff")
_EMPTY_EVENTS = pl.DataFrame(schema={"vt_symbol": pl.Utf8, "eff": pl.Datetime("us")})
def _load_extra_domains(codes: list[str], data_dir: str,
start: str, end: str, out_cols: list[str]) -> dict:
"""P1-B 四新域按引用列加载(引用瘦身的域级延伸: 未引用的域不读盘)."""
cols = set(out_cols)
return {
"dividend": (_load_dividend_cum(codes, data_dir)
if "send_total_12m" in cols else _EMPTY_EVENTS),
"gdhs": (_load_gdhs_events(codes, data_dir)
if "gdhs_chg" in cols else _EMPTY_EVENTS),
"topholder": (_load_topholder_events(codes, data_dir)
if "topholder_chg" in cols else _EMPTY_EVENTS),
"vb": (_load_vb_daily(codes, data_dir, start, end)
if "ep_vb" in cols or "bp_vb" in cols else _EMPTY_EVENTS),
}
# ==================== 对外主入口 ====================
# 分块大小: NAS 7.9G 物理内存下,全市场 grid(5555股×2670日×33列≈4G)单次
@@ -680,7 +1031,8 @@ def _load_feature_events(codes: list[str], data_dir: str) -> tuple[pl.DataFrame,
"""报告期特征事件 + forecast 事件 + 兑现差事件(全 codes 一次加载,
分块 join 共用右表)."""
stmt_cols = [c for c in FEATURE_COLUMNS
if c not in _FORECAST_COLS and c not in _BEAT_COLS]
if c not in _FORECAST_COLS and c not in _BEAT_COLS
and c not in _DOMAIN_COLS]
reports = _load_statements(codes, data_dir)
if reports.height == 0:
stmt_events = pl.DataFrame(schema={
@@ -754,6 +1106,7 @@ def iter_fundamental_feature_chunks(
"""
out_cols = list(columns) if columns is not None else list(FEATURE_COLUMNS)
stmt_events, fc_events, beat_events = _load_feature_events(codes, data_dir)
extra = _load_extra_domains(codes, data_dir, start, end, out_cols)
day_list = _prepare_days(start, end, trading_dates)
for i in range(0, len(codes), batch_codes):
chunk_codes = codes[i:i + batch_codes]
@@ -773,6 +1126,41 @@ def iter_fundamental_feature_chunks(
out = out.sort("datetime").join_asof(
beat_events, left_on="datetime", right_on="eff",
by="vt_symbol", strategy="backward").drop("eff")
# ---- P1-B 四新域 ----
div = extra["dividend"]
if div.height:
# E11: cum(d) cum(d365) = 窗 (d365, d] 事件合计(开区间下界);
# 域内无事件(≤d) → NaN 不填 0(与缺文件/域不覆盖区分)
div365 = div.rename({"_send_cum": "_send_cum_365"})
out = out.sort("datetime").join_asof(
div, left_on="datetime", right_on="eff",
by="vt_symbol", strategy="backward").drop("eff")
out = (
out.with_columns(
(pl.col("datetime") - pl.duration(days=365)).alias("_dt365"))
.sort("_dt365")
.join_asof(div365, left_on="_dt365", right_on="eff",
by="vt_symbol", strategy="backward")
.drop("eff", "_dt365")
.with_columns(
pl.when(pl.col("_send_cum").is_null()).then(None)
.otherwise(pl.col("_send_cum")
- pl.col("_send_cum_365").fill_null(0.0))
.alias("send_total_12m")))
for key in ("gdhs", "topholder"):
if extra[key].height:
out = out.sort("datetime").join_asof(
extra[key], left_on="datetime", right_on="eff",
by="vt_symbol", strategy="backward").drop("eff")
if extra["vb"].height:
out = out.sort("datetime").join_asof(
extra["vb"], left_on="datetime", right_on="eff",
by="vt_symbol", strategy="backward").drop("eff")
# 缺域/无事件兜底: 输出列缺失 → null 列(优雅降级,select 不炸)
absent = [c for c in out_cols if c not in out.columns]
if absent:
out = out.with_columns(
[pl.lit(None, dtype=pl.Float64).alias(c) for c in absent])
yield out.sort(["vt_symbol", "datetime"]).select(
["vt_symbol", "datetime", *out_cols])
grid = out = None # 批间释放(下一批重绑定)
+151 -1
View File
@@ -12,7 +12,7 @@ if _VNPY_SRC not in sys.path:
# 6 只股 × 24 报告期(2018Q1~2023Q4;2018-2019 为 16 季滚动窗/5 年 CAGR 预热史),
# NOTICE_DATE 错位: Q1→当年4-28 / H1→当年8-29 / Q3→当年10-27 / 年报→次年4-25
# 数值全部手算可验证(断言用),公式见各 build 函数内注释。
from datetime import date
from datetime import date, datetime
import polars as pl
@@ -33,6 +33,7 @@ NP_Q = ([7, 10, 11, 12] + [8, 11, 12, 13] # 2018 / 2019
TA_OF = lambda i: 1000 + 50 * max(i - 12, 0)
EQ_OF = lambda i: 500 + 20 * max(i - 12, 0)
AR_OF = lambda i: 100 + 10 * max(i - 12, 0)
AP_OF = lambda i: 120 + 8 * max(i - 12, 0) # P1-B: 应付账款(E08 CCC 分子)
SC_OF = lambda i: 100 if i < 20 else 110 # 2023 起股本扩张 10%(NSI=0.1)
SYN_STOCKS = {
@@ -119,6 +120,8 @@ def build_synthetic_static(root: str) -> str:
"SHORT_LOAN": s * (100 + max(i - 12, 0)),
# P1 新原料: 存货/货币资金(银行无存货 → null 传播)
"INVENTORY": None if is_bank else s * (200 + 5 * max(i - 12, 0)),
# P1-B: 应付账款(银行无应付账款概念 → null 传播,CCC 剔金融)
"ACCOUNTS_PAYABLE": None if is_bank else s * AP_OF(i),
"MONETARYFUNDS": s * (150 + 10 * max(i - 12, 0))})
cashflow_rows.append({**common,
"NETCASH_OPERATE": 1.2 * np_c, "SALES_SERVICES": 1.05 * rev_c,
@@ -160,6 +163,153 @@ def build_synthetic_static(root: str) -> str:
{"股票代码": "600000", "预测指标": "归属于上市公司股东的净利润",
"预测数值": 55.0, "业绩变动幅度": 15.0, "预告类型": "预增", "公告日期": date(2023, 5, 10)},
]).write_parquet(os.path.join(static_dir, "forecast", "20221231_forecast.parquet"))
# ==================== P1-B 批四新域(dividend/gdhs/top_holders/vb) ====================
# 契约=NAS 实测 2026-09-08: dividend 日期/数值列为字符串('' 为空);
# gdhs 截止日列名「股东户数统计截止日-本次」(无第二连字符);vb date 为字符串。
# ---- dividend 按股文件(A/B 两只;C/BANK 无文件 → send_total_12m NaN)----
# A 事件(eff = 回退链首个非空公告日):
# 2022-06-15 送转 0.2 / 2023-03-10 送 0.5 / 2023-09-20 送 0.3
# 2023-12-01(仅 dividPreNoticeDate,回退链用例) 送 0.4
# 全链空行(dividStocksPs=9.9 哨兵,应被丢弃)
# 2024-01-05 送 1.0(迟到事件,2024 起才可见)
os.makedirs(os.path.join(static_dir, "dividend"), exist_ok=True)
pl.DataFrame([
{"code": "sh.600000", "dividPreNoticeDate": "", "dividAgmPumDate": "",
"dividPlanAnnounceDate": "2022-06-15", "dividPlanDate": "",
"dividStocksPs": "0.200000", "dividCashPsBeforeTax": "1.000000"},
{"code": "sh.600000", "dividPreNoticeDate": "", "dividAgmPumDate": "",
"dividPlanAnnounceDate": "2023-03-10", "dividPlanDate": "",
"dividStocksPs": "0.500000", "dividCashPsBeforeTax": "2.000000"},
{"code": "sh.600000", "dividPreNoticeDate": "", "dividAgmPumDate": "",
"dividPlanAnnounceDate": "2023-09-20", "dividPlanDate": "",
"dividStocksPs": "0.300000", "dividCashPsBeforeTax": "3.000000"},
{"code": "sh.600000", "dividPreNoticeDate": "2023-12-01", "dividAgmPumDate": "",
"dividPlanAnnounceDate": "", "dividPlanDate": "",
"dividStocksPs": "0.400000", "dividCashPsBeforeTax": "4.000000"},
{"code": "sh.600000", "dividPreNoticeDate": "", "dividAgmPumDate": "",
"dividPlanAnnounceDate": "", "dividPlanDate": "",
"dividStocksPs": "9.900000", "dividCashPsBeforeTax": "9.900000"},
{"code": "sh.600000", "dividPreNoticeDate": "", "dividAgmPumDate": "",
"dividPlanAnnounceDate": "2024-01-05", "dividPlanDate": "",
"dividStocksPs": "1.000000", "dividCashPsBeforeTax": "5.000000"},
]).write_parquet(os.path.join(static_dir, "dividend", "600000.SH_dividend.parquet"))
# B: 纯现金分红事件(dividStocksPs=0 → send_total_12m = 0 非 NaN)
pl.DataFrame([
{"code": "sz.000001", "dividPreNoticeDate": "", "dividAgmPumDate": "",
"dividPlanAnnounceDate": "2023-03-10", "dividPlanDate": "",
"dividStocksPs": "0.000000", "dividCashPsBeforeTax": "0.500000"},
]).write_parquet(os.path.join(static_dir, "dividend", "000001.SZ_dividend.parquet"))
# ---- gdhs 按期文件(事件表: 截止日→户数,PIT=公告日期)----
os.makedirs(os.path.join(static_dir, "gdhs"), exist_ok=True)
def _gdhs_rows(items):
return [{"代码": c, "名称": f"{c}", "股东户数-本次": h,
"股东户数统计截止日-本次": cutoff, "公告日期": ann,
"股东户数-增减比例": -1.0} # 噪声列(红线: 不可直接用,应无影响)
for c, h, cutoff, ann in items]
pl.DataFrame(_gdhs_rows([
("000001", 46000, date(2022, 6, 30), date(2022, 7, 20)),
])).write_parquet(os.path.join(static_dir, "gdhs", "20220630_gdhs.parquet"))
pl.DataFrame(_gdhs_rows([
("600000", 84000, date(2022, 9, 30), date(2022, 10, 20)), # A 首事件(Δ=null)
("000001", 48000, date(2022, 9, 30), date(2022, 10, 20)),
])).write_parquet(os.path.join(static_dir, "gdhs", "20220930_gdhs.parquet"))
pl.DataFrame(_gdhs_rows([
("600000", 80000, date(2022, 12, 31), date(2023, 1, 10)),
("000001", 50000, date(2022, 12, 31), date(2023, 1, 10)),
])).write_parquet(os.path.join(static_dir, "gdhs", "20221231_gdhs.parquet"))
pl.DataFrame(_gdhs_rows([
("600000", 72000, date(2023, 3, 31), date(2023, 4, 20)),
])).write_parquet(os.path.join(static_dir, "gdhs", "20230331_gdhs.parquet"))
# 同 (股, 截止日) 两行不同公告日 → 去重取最新公告(08-01 修订取代 07-20)
pl.DataFrame(_gdhs_rows([
("600000", 68000, date(2023, 6, 30), date(2023, 7, 25)),
("000001", 50000, date(2023, 6, 30), date(2023, 7, 20)),
]) + _gdhs_rows([
("000001", 50500, date(2023, 6, 30), date(2023, 8, 1)),
])).write_parquet(os.path.join(static_dir, "gdhs", "20230630_gdhs.parquet"))
pl.DataFrame(_gdhs_rows([
("600000", 77760, date(2023, 9, 30), date(2023, 10, 15)),
])).write_parquet(os.path.join(static_dir, "gdhs", "20230930_gdhs.parquet"))
# ---- top_holders: 原始按股×期小文件树(供预聚合脚本测试)+ 聚合产物 ----
# 600000 五期占比合计: 50/57/60/65/62;000001 单期 40(单期 → Δ=null)
th_dir = os.path.join(static_dir, "top_holders")
os.makedirs(th_dir, exist_ok=True)
_TH_SUMS = {("600000.SH", "20220630"): 5.0, ("600000.SH", "20220930"): 5.7,
("600000.SH", "20221231"): 6.0, ("600000.SH", "20230331"): 6.5,
("600000.SH", "20230630"): 6.2, ("000001.SZ", "20221231"): 4.0}
for (file_code, period), per_row in _TH_SUMS.items():
rows = []
for r in range(10):
rows.append({
"名次": r + 1, "股东名称": f"股东{r}", "股东性质": "基金",
"股份类型": "流通A股", "持股数": 1000 + r,
"占总流通股本持股比例": (None if (file_code == "600000.SH" and r == 9)
else per_row),
# 增减混合类型(「不变」字符串+数值字符串,coalesce/求和不受影响)
"增减": "不变" if r % 2 else f"{1000 + r}", "变动比率": 0.0,
})
pl.DataFrame(rows).write_parquet(
os.path.join(th_dir, f"{file_code}_{period}_top_holders.parquet"))
# 聚合产物(adapter 只读这个;与脚本对同一小树的输出逐值一致)
root = os.path.dirname(static_dir)
agg_dir = os.path.join(root, "factor_cache")
os.makedirs(agg_dir, exist_ok=True)
pl.DataFrame([
{"file_code": fc, "period": datetime.strptime(p, "%Y%m%d").date(),
"hold_pct": 10.0 * v, "n_holders": 10}
for (fc, p), v in _TH_SUMS.items()
]).write_parquet(os.path.join(agg_dir, "top_holders_agg.parquet"))
# ---- valuation_baostock 按年文件(与 static 同级的兄弟目录)----
vb_dir = os.path.join(root, "valuation_baostock")
os.makedirs(vb_dir, exist_ok=True)
pl.DataFrame([
{"symbol": "600000", "exchange": "SH", "date": "2023-08-25",
"peTTM": 12.0, "psTTM": 6.0, "pcfNcfTTM": 8.0, "pbMRQ": 2.4,
"turn": 0.5, "pctChg": 1.0, "isST": 0},
{"symbol": "600000", "exchange": "SH", "date": "2023-08-28",
"peTTM": 11.0, "psTTM": 6.0, "pcfNcfTTM": 8.0, "pbMRQ": 2.2,
"turn": 0.5, "pctChg": 1.0, "isST": 0},
{"symbol": "600000", "exchange": "SH", "date": "2023-08-29",
"peTTM": 10.0, "psTTM": 6.0, "pcfNcfTTM": 8.0, "pbMRQ": 2.0,
"turn": 0.5, "pctChg": 1.0, "isST": 0},
{"symbol": "000001", "exchange": "SZ", "date": "2023-08-29",
"peTTM": 20.0, "psTTM": 6.0, "pcfNcfTTM": 8.0, "pbMRQ": 0.8,
"turn": 0.5, "pctChg": 1.0, "isST": 0},
# 亏损股: baostock 不给负 pe → NaN;pe=0 守卫用例
{"symbol": "300001", "exchange": "SZ", "date": "2023-08-29",
"peTTM": None, "psTTM": 6.0, "pcfNcfTTM": 8.0, "pbMRQ": 5.0,
"turn": 0.5, "pctChg": 1.0, "isST": 0},
{"symbol": "000001", "exchange": "SZ", "date": "2023-08-30",
"peTTM": 0.0, "psTTM": 6.0, "pcfNcfTTM": 8.0, "pbMRQ": 0.8,
"turn": 0.5, "pctChg": 1.0, "isST": 0},
# exchange 全称变体(SSE/SZSE 旧契约形态)一行 → 映射兼容锁定
{"symbol": "600004", "exchange": "SSE", "date": "2023-08-29",
"peTTM": 40.0, "psTTM": 6.0, "pcfNcfTTM": 8.0, "pbMRQ": 4.0,
"turn": 0.5, "pctChg": 1.0, "isST": 0},
]).write_parquet(os.path.join(vb_dir, "2023.parquet"))
# 2026 文件: 08-13 起(bs 日喂起点;01-01~08-12 缺口由 static/valuation 补)
pl.DataFrame([
{"symbol": "600000", "exchange": "SH", "date": "2026-08-13",
"peTTM": 20.0, "psTTM": 6.0, "pcfNcfTTM": 8.0, "pbMRQ": 4.0,
"turn": 0.5, "pctChg": 1.0, "isST": 0},
{"symbol": "600000", "exchange": "SH", "date": "2026-08-14",
"peTTM": 19.0, "psTTM": 6.0, "pcfNcfTTM": 8.0, "pbMRQ": 3.9,
"turn": 0.5, "pctChg": 1.0, "isST": 0},
]).write_parquet(os.path.join(vb_dir, "2026.parquet"))
# ---- static/valuation(ak em 中文列;仅补 vb 2026-01-01~08-12 缺口)----
os.makedirs(os.path.join(static_dir, "valuation"), exist_ok=True)
pl.DataFrame([
{"数据日期": date(2026, 1, 5), "PE(TTM)": 25.0, "市净率": 5.0},
{"数据日期": date(2026, 8, 12), "PE(TTM)": 24.0, "市净率": 4.8},
]).write_parquet(os.path.join(static_dir, "valuation", "600000.SH_valuation.parquet"))
return static_dir
@@ -0,0 +1,248 @@
# tests/factor/test_fundamental_p1b_adapter.py
"""P1-B 批财务因子适配层: 四新域(dividend/gdhs/top_holders agg/vb)+ CCC 三表自算.
口径锚(P1-B 任务书 + NAS 实测 2026-09-08):
- E08 CCC = 365×(AR_avg/REV_TTM + INV_avg/COGS_TTM AP_avg/COGS_TTM);
均值 = t t4 报告期期末余额平均(恰隔 4 季守卫);金融股置 NaN
- E11 send_total_12m = 365 (开区间下界)分红事件 dividStocksPs 累计和;
PIT 日期回退链 dividPlanAnnounceDatePreNoticeAgmPumPlan,全空行丢弃;
域内无事件(决策日) NaN(不填 0),纯现金事件送转 0 0
- E12 gdhs_chg = 相邻事件 Δln(户数)(按截止日排序,公告日期 asof;
(,截止日) 多行取最新公告)
- E13 topholder_chg = 相邻期十大流通合计占比 Δpct;PIT = 法定披露截止
(Q104-30/H108-31/Q310-31/年报次年04-30,保守侧)
- D14 ep_vb/bp_vb = 1/peTTM1/pbMRQ(0 NaN);2026-01-01~08-12 缺口由
static/valuation(ak 中文列)倒数补,08-13 起用 vb
"""
import math
import shutil
import sys, os
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..")))
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..", "vnpy_v4.4.0")))
import pytest
from sanguo_factor.fundamental_adapter import build_fundamental_features
A, B, C = "600000.SSE", "000001.SZSE", "300001.SZSE"
BANK = "601398.SSE"
# 合成 dividend 事件(eff, dividStocksPs;A)
_A_EVENTS = [("2022-06-15", 0.2), ("2023-03-10", 0.5), ("2023-09-20", 0.3),
("2023-12-01", 0.4), ("2024-01-05", 1.0)]
def _dt(day: str):
import datetime as _d
return _d.datetime.strptime(day, "%Y-%m-%d")
def val(df, vt: str, day: str, col: str):
row = df.filter((df["vt_symbol"] == vt) & (df["datetime"] == _dt(day)))
assert row.height == 1, f"grid 缺行 {vt} {day}"
v = row[col][0]
return None if v is None else float(v)
def _window_sum(day: str) -> float:
"""独立重算: (day-365, day] 窗内 A 事件 dividStocksPs 合计."""
import datetime as _d
d = _dt(day)
lo = d.replace(year=d.year - 1)
return sum(v for eff, v in _A_EVENTS if _dt(eff) <= d and _dt(eff) > lo)
@pytest.fixture(scope="module")
def feat(synthetic_static):
return build_fundamental_features(
[A, B, C], "2023-01-01", "2023-12-31", data_dir=synthetic_static)
# ---------- E08 CCC(三表自算) ----------
def test_ccc_value_at_h1(feat):
# 2023H1(i=21): AR mean(190,150)=170 INV mean(245,225)=235 AP mean(192,176)=176
# REV_TTM=600 COGS_TTM=0.6×600=360
expect = 365 * (170 / 600 + 235 / 360 - 176 / 360)
assert val(feat, A, "2023-08-29", "ccc") == pytest.approx(expect)
# scale 不变性(C=2×): 三均值与两 TTM 同比放大 → CCC 不变
assert val(feat, C, "2023-08-29", "ccc") == pytest.approx(expect)
def test_ccc_pit_boundary(feat):
# 2023-08-28 见 2023Q1(i=20): AR mean(180,140)=160 INV mean(240,220)=230
# AP mean(184,152)=168 REV_TTM=q(17..20)=600 COGS_TTM=360
assert val(feat, A, "2023-08-28", "ccc") == pytest.approx(
365 * (160 / 600 + 230 / 360 - 168 / 360))
assert val(feat, A, "2023-08-29", "ccc") == pytest.approx(
365 * (170 / 600 + 235 / 360 - 176 / 360))
def test_ccc_guards(feat, synthetic_static):
# B 缺 2022Q1(i=16) → 2023Q1 均值基期断档 → NaN;H1 基期(i=17) 恢复
assert val(feat, B, "2023-04-28", "ccc") is None
assert val(feat, B, "2023-08-29", "ccc") == pytest.approx(
365 * (170 / 600 + 235 / 360 - 176 / 360))
# 金融股: CCC 属营运效率,银行无存货/应付概念 → NaN(§7 红线扩展)
df = build_fundamental_features(
[BANK], "2023-08-25", "2023-09-02", data_dir=synthetic_static)
assert val(df, BANK, "2023-08-29", "ccc") is None
# ---------- E11 高送转强度(dividend 域) ----------
def test_send_total_values_and_fallback_chain(feat):
# 回退链: 2023-12-01 事件只有 dividPreNoticeDate → 计入
assert val(feat, A, "2023-11-30", "send_total_12m") == pytest.approx(0.8)
assert val(feat, A, "2023-12-01", "send_total_12m") == pytest.approx(1.2)
# 全链空行(dividStocksPs=9.9 哨兵)被丢弃: 任何时点 ≤2.2
assert val(feat, A, "2023-12-31", "send_total_12m") == pytest.approx(1.2)
def test_send_total_pit_window_boundaries(synthetic_static):
df = build_fundamental_features(
[A], "2023-03-01", "2024-03-15", data_dir=synthetic_static)
# 事件日边界(当日可见)
assert val(df, A, "2023-03-09", "send_total_12m") == pytest.approx(0.2)
assert val(df, A, "2023-03-10", "send_total_12m") == pytest.approx(0.7)
# 365 天滑出边界(开区间下界): 2022-06-15 事件在 2023-06-14 仍在/06-15 出窗
assert val(df, A, "2023-06-14", "send_total_12m") == pytest.approx(0.7)
assert val(df, A, "2023-06-15", "send_total_12m") == pytest.approx(0.5)
# 迟到事件: 2024-01-05 公告,04 日不可见/05 日计入
assert val(df, A, "2024-01-04", "send_total_12m") == pytest.approx(1.2)
assert val(df, A, "2024-01-05", "send_total_12m") == pytest.approx(2.2)
# 365 天出窗边界 = 精确 365×24h(跨 2024 闰日 → 2023-03-10 事件在
# 2024-03-09(=事件+366 日历日)已出窗)
assert val(df, A, "2024-03-08", "send_total_12m") == pytest.approx(2.2)
assert val(df, A, "2024-03-09", "send_total_12m") == pytest.approx(1.7)
def test_send_total_domain_semantics(feat):
# B: 域内纯现金事件(dividStocksPs=0)→ 0(非 NaN);事件前无事件 → NaN
assert val(feat, B, "2023-03-09", "send_total_12m") is None
assert val(feat, B, "2023-03-10", "send_total_12m") == pytest.approx(0.0)
# C: 无 dividend 文件(域不覆盖)→ NaN(不填 0)
assert val(feat, C, "2023-12-31", "send_total_12m") is None
# ---------- E12 股东户数变化(gdhs 域) ----------
def test_gdhs_chg_adjacent_events_and_pit(feat):
ln = math.log
# A 首事件(2022-09-30)Δ=null;2022-12-31 事件公告 2023-01-10
assert val(feat, A, "2023-01-09", "gdhs_chg") is None
assert val(feat, A, "2023-01-10", "gdhs_chg") == pytest.approx(ln(80000 / 84000))
# 相邻事件各自 Δln(Q1→Q2 换值,非全史一个值): 公告日边界
assert val(feat, A, "2023-04-19", "gdhs_chg") == pytest.approx(ln(80000 / 84000))
assert val(feat, A, "2023-04-20", "gdhs_chg") == pytest.approx(ln(72000 / 80000))
assert val(feat, A, "2023-07-24", "gdhs_chg") == pytest.approx(ln(72000 / 80000))
assert val(feat, A, "2023-07-25", "gdhs_chg") == pytest.approx(ln(68000 / 72000))
assert val(feat, A, "2023-10-14", "gdhs_chg") == pytest.approx(ln(68000 / 72000))
assert val(feat, A, "2023-10-15", "gdhs_chg") == pytest.approx(ln(77760 / 68000))
def test_gdhs_chg_dedupe_latest_announcement(feat):
# B 2023-06-30 截止日两行(公告 07-20 户数 50000 / 08-01 修订 50500):
# 去重取最新公告 → 08-01 前不得提前泄露修订值
assert val(feat, B, "2023-07-31", "gdhs_chg") == pytest.approx(math.log(50000 / 48000))
assert val(feat, B, "2023-08-01", "gdhs_chg") == pytest.approx(math.log(50500 / 50000))
def test_gdhs_chg_uncovered_stock(feat):
# C 无 gdhs 行 → NaN;「增减比例」噪声列不参与(值由 Δln 自算)
assert val(feat, C, "2023-12-31", "gdhs_chg") is None
# ---------- E13 十大流通股东占比变化(top_holders 聚合域) ----------
def test_topholder_chg_statutory_deadlines(feat):
# 法定披露截止 PIT: Q1→04-30 / 年报→次年04-30 / H1→08-31 / Q3→10-31
# Q1 与年报共用 04-30 截止 → 同日双事件 asof 取期更近的 Q1 事件
assert val(feat, A, "2023-04-29", "topholder_chg") == pytest.approx(7.0) # Q3'22: 5750
assert val(feat, A, "2023-04-30", "topholder_chg") == pytest.approx(5.0) # Q1'23: 6560(年报 Δ=3 同日让位)
assert val(feat, A, "2023-08-30", "topholder_chg") == pytest.approx(5.0)
assert val(feat, A, "2023-08-31", "topholder_chg") == pytest.approx(-3.0) # H1'23: 6265
# 无 Q3'23 期数据 → 10-31 无法定截止新事件,值持续
assert val(feat, A, "2023-12-31", "topholder_chg") == pytest.approx(-3.0)
def test_topholder_chg_insufficient_history(feat):
# B 仅单期 → 无上期 → NaN
assert val(feat, B, "2023-12-31", "topholder_chg") is None
# C 不在聚合产物 → NaN
assert val(feat, C, "2023-12-31", "topholder_chg") is None
# ---------- D14 长史 EP/BP(vb 域 + 2026 补口) ----------
def test_vb_ep_bp_values_and_ffill(feat):
# A: 08-25 pe=12 → 周末 08-26(日历 grid)前向填充;08-28/08-29 换值
assert val(feat, A, "2023-08-26", "ep_vb") == pytest.approx(1 / 12)
assert val(feat, A, "2023-08-28", "ep_vb") == pytest.approx(1 / 11)
assert val(feat, A, "2023-08-29", "ep_vb") == pytest.approx(1 / 10)
assert val(feat, A, "2023-08-29", "bp_vb") == pytest.approx(1 / 2.0)
assert val(feat, B, "2023-08-29", "ep_vb") == pytest.approx(1 / 20)
assert val(feat, B, "2023-08-29", "bp_vb") == pytest.approx(1 / 0.8)
def test_vb_ep_guards(feat):
# 亏损(baostock pe=NaN)→ ep NaN;pe=0 → NaN(1/0 守卫);bp 不受影响
assert val(feat, C, "2023-08-29", "ep_vb") is None
assert val(feat, C, "2023-08-29", "bp_vb") == pytest.approx(1 / 5.0)
assert val(feat, B, "2023-08-30", "ep_vb") is None
assert val(feat, B, "2023-08-30", "bp_vb") == pytest.approx(1 / 0.8)
def test_vb_exchange_full_form_tolerated(synthetic_static):
"""exchange 全称变体(SSE/SZSE)与实测缩写(SH/SZ)同被映射(旧契约兼容)."""
df = build_fundamental_features(
["600004.SSE", "600000.SSE"], "2023-08-25", "2023-08-31",
data_dir=synthetic_static)
assert val(df, "600004.SSE", "2023-08-29", "ep_vb") == pytest.approx(1 / 40)
assert val(df, "600004.SSE", "2023-08-29", "bp_vb") == pytest.approx(1 / 4.0)
def test_vb_2026_gap_splice(synthetic_static):
"""2026-01-01~08-12 用 static/valuation(ak)倒数补;08-13 起用 vb."""
df = build_fundamental_features(
[A], "2026-01-01", "2026-08-15", data_dir=synthetic_static)
assert val(df, A, "2026-01-05", "ep_vb") == pytest.approx(1 / 25) # ak
assert val(df, A, "2026-01-05", "bp_vb") == pytest.approx(1 / 5.0)
assert val(df, A, "2026-08-11", "ep_vb") == pytest.approx(1 / 25) # ak ffill
assert val(df, A, "2026-08-12", "ep_vb") == pytest.approx(1 / 24) # ak 末日
assert val(df, A, "2026-08-13", "ep_vb") == pytest.approx(1 / 20) # vb 接管
assert val(df, A, "2026-08-13", "bp_vb") == pytest.approx(1 / 4.0)
assert val(df, A, "2026-08-14", "ep_vb") == pytest.approx(1 / 19)
# ---------- 缺域优雅降级(本地无 NAS 数据) ----------
def test_missing_domains_degrade_to_nan(synthetic_static, tmp_path):
"""四域整体缺失 → 六新列全 NaN + warning,不崩(管线可本地跑通)."""
dst = tmp_path / "no_domains"
shutil.copytree(os.path.dirname(synthetic_static), str(dst))
for sub in ("dividend", "gdhs", "valuation"):
shutil.rmtree(str(dst / "static" / sub))
shutil.rmtree(str(dst / "factor_cache"))
shutil.rmtree(str(dst / "valuation_baostock"))
with pytest.warns(UserWarning, match=""):
df = build_fundamental_features(
[A], "2023-08-25", "2023-08-31", data_dir=str(dst / "static"))
assert df.height > 0
for col in ("send_total_12m", "gdhs_chg", "topholder_chg", "ep_vb", "bp_vb"):
assert df[col].null_count() == df.height, f"{col} 缺域应全 NaN"
# 报表域仍在 → ccc 照常产出(缺的只是四个新域)
assert val(df, A, "2023-08-29", "ccc") is not None
# ---------- 分块等值(新列过一遍 batch_codes=1 极端路径) ----------
def test_chunked_equals_full_p1b_columns(synthetic_static):
six = ["600000.SSE", "000001.SZSE", "300001.SZSE",
"600004.SSE", "000333.SZSE", "300124.SZSE"]
full = build_fundamental_features(
six, "2023-01-01", "2023-12-31", data_dir=synthetic_static, batch_codes=6)
by_one = build_fundamental_features(
six, "2023-01-01", "2023-12-31", data_dir=synthetic_static, batch_codes=1)
key = ["vt_symbol", "datetime"]
assert full.sort(key).equals(by_one.sort(key))