diff --git a/sanguo_factor/fundamental_adapter.py b/sanguo_factor/fundamental_adapter.py index 6b68f3e..078caf8 100644 --- a/sanguo_factor/fundamental_adapter.py +++ b/sanguo_factor/fundamental_adapter.py @@ -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+DIO−DPO(缺任一原料 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(d−365) = 窗 (d−365, 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(d−365) = 窗 (d−365, 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 # 批间释放(下一批重绑定) diff --git a/tests/factor/conftest.py b/tests/factor/conftest.py index 709fd5b..b4ac07f 100644 --- a/tests/factor/conftest.py +++ b/tests/factor/conftest.py @@ -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 diff --git a/tests/factor/test_fundamental_p1b_adapter.py b/tests/factor/test_fundamental_p1b_adapter.py new file mode 100644 index 0000000..a6e5815 --- /dev/null +++ b/tests/factor/test_fundamental_p1b_adapter.py @@ -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 与 t−4 报告期期末余额平均(恰隔 4 季守卫);金融股置 NaN +- E11 send_total_12m = 近 365 天(开区间下界)分红事件 dividStocksPs 累计和; + PIT 日期回退链 dividPlanAnnounceDate→PreNotice→AgmPum→Plan,全空行丢弃; + 域内无事件(≤决策日) → NaN(不填 0),纯现金事件送转 0 → 0 +- E12 gdhs_chg = 相邻事件 Δln(户数)(按截止日排序,公告日期 asof; + 同 (股,截止日) 多行取最新公告) +- E13 topholder_chg = 相邻期十大流通合计占比 Δpct;PIT = 法定披露截止 + (Q1→04-30/H1→08-31/Q3→10-31/年报→次年04-30,保守侧) +- D14 ep_vb/bp_vb = 1/peTTM、1/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: 57−50 + assert val(feat, A, "2023-04-30", "topholder_chg") == pytest.approx(5.0) # Q1'23: 65−60(年报 Δ=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: 62−65 + # 无 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))