From 98c3339f8a0e5d837205d13d50500ce53a9abe3d Mon Sep 17 00:00:00 2001 From: claude_dev Date: Thu, 10 Sep 2026 07:18:28 +0800 Subject: [PATCH] =?UTF-8?q?refactor(factor):=20fundamental=5Fadapter=20120?= =?UTF-8?q?1=E8=A1=8C=E6=8B=86=E4=B8=83=E6=A8=A1=E5=9D=97=E2=80=94?= =?UTF-8?q?=E2=80=94=E9=97=A8=E9=9D=A2re-export=E9=9B=B6=E8=A1=8C=E4=B8=BA?= =?UTF-8?q?=E5=8F=98=E5=8C=96,AST=E9=80=90=E5=AE=9A=E4=B9=8952/52=E5=AD=97?= =?UTF-8?q?=E8=8A=82=E4=B8=80=E8=87=B4+1143=E6=B5=8B=E8=AF=95=E4=B8=8D?= =?UTF-8?q?=E6=94=B9=E4=B8=80=E5=AD=97=E5=85=A8=E7=BB=BF=20[nas]?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- sanguo_factor/fundamental_adapter.py | 1235 ++---------------- sanguo_factor/fundamental_domains.py | 329 +++++ sanguo_factor/fundamental_forecast.py | 118 ++ sanguo_factor/fundamental_pit.py | 223 ++++ sanguo_factor/fundamental_report_features.py | 313 +++++ sanguo_factor/fundamental_schema.py | 128 ++ sanguo_factor/fundamental_statements.py | 130 ++ 7 files changed, 1314 insertions(+), 1162 deletions(-) create mode 100644 sanguo_factor/fundamental_domains.py create mode 100644 sanguo_factor/fundamental_forecast.py create mode 100644 sanguo_factor/fundamental_pit.py create mode 100644 sanguo_factor/fundamental_report_features.py create mode 100644 sanguo_factor/fundamental_schema.py create mode 100644 sanguo_factor/fundamental_statements.py diff --git a/sanguo_factor/fundamental_adapter.py b/sanguo_factor/fundamental_adapter.py index 4bd7f27..c91a15c 100644 --- a/sanguo_factor/fundamental_adapter.py +++ b/sanguo_factor/fundamental_adapter.py @@ -34,1168 +34,79 @@ P1-B 四新域(NAS 实测 2026-09-08): 输出: vt_symbol × datetime(日频) × FEATURE_COLUMNS,join 到 alpha_df 后 供 cs_rank(列) 表达式直接消费。 + +模块拆分(2026-09-10,纯结构重构零行为变化): 本文件保留为门面,re-export +全部原模块级名字(引用方与测试零改动),实现按数据域拆为: +- fundamental_schema.py 常量/列契约(FEATURE_COLUMNS 等) +- fundamental_statements.py 三表读取层(报告期宽表) +- fundamental_report_features.py 报告期级指标预计算(单季化/TTM/比率族) +- fundamental_forecast.py 业绩预告事件层 +- fundamental_domains.py P1-B 四新域事件层 +- fundamental_pit.py PIT 对齐 + 分块 join 主入口 """ from __future__ import annotations -import os -import warnings -from datetime import date, datetime - -import polars as pl - -DEFAULT_STATIC_DIR = "/volume1/stock/sanguo_vnpy_v2/data/static" -_SYM = "vt_symbol" # 时序分组键 - -# vt_symbol 后缀 → NAS 文件名后缀(600000.SSE → 600000.SH) -_VT_TO_FILE_SUFFIX = {"SSE": "SH", "SZSE": "SZ"} - -# forecast 预告类型 → 有序分(§3.6 F04) -FORECAST_TYPE_SCORE: dict[str, int] = { - "预增": 3, "略增": 2, "扭亏": 2, "续盈": 1, "减亏": 1, "不确定": 0, - "略减": -1, "增亏": -2, "续亏": -2, "预减": -3, "首亏": -3, -} - -# 累计口径列映射: 源列 → 短名(需单季化+TTM) -_CUM_MAP = { - "TOTAL_OPERATE_INCOME": "rev", - "OPERATE_COST": "cogs", - "PARENT_NETPROFIT": "np", - "DEDUCT_PARENT_NETPROFIT": "dnp", - "TOTAL_PROFIT": "tp", - "INVEST_INCOME": "invest_inc", - "FAIRVALUE_CHANGE_INCOME": "fv_inc", - "ASSET_IMPAIRMENT_LOSS": "asset_imp", - "CREDIT_IMPAIRMENT_LOSS": "credit_imp", - "NETCASH_OPERATE": "cfo", - "SALES_SERVICES": "sales_cash", - "ACCEPT_INVEST_CASH": "acc_inv_cash", - # ---- P1 批新增(income;NAS 实测 2026-09-08 列名核实)---- - "RESEARCH_EXPENSE": "research", # 研发费用(2018Q3 起单列,早年 null 传播) - "SALE_EXPENSE": "sale_exp", # 销售费用 - "FE_INTEREST_EXPENSE": "int_exp", # 财务费用-利息费用(披露稀疏,EBIT 组装层按 0) - "INCOME_TAX": "tax", # 所得税费用(τ 实际税率分子) - "BASIC_EPS": "eps", # 基本每股收益(累计口径,需单季化) - # ---- P1 批新增(cashflow 补充资料段 + 融资流)---- - "FA_IR_DEPR": "fa_depr", # 固定资产折旧 - "IA_AMORTIZE": "ia_amort", # 无形资产摊销 - "LPE_AMORTIZE": "lpe_amort", # 长期待摊费用摊销 - "USERIGHT_ASSET_AMORTIZE": "ua_amort", # 使用权资产折旧摊销(2019 起才有) - "CONSTRUCT_LONG_ASSET": "capex", # 购建长期资产现金(C17 投资率/D15 FCF) - "RECEIVE_LOAN_CASH": "recv_loan", # 取得借款现金(E03) - "ISSUE_BOND": "issue_bond", # 发行债券现金(E03) - "PAY_DEBT_CASH": "pay_debt", # 偿还债务现金(E03) -} -# 存量列(balance,时点值直接用;缺列 → null) -_BALANCE_COLS = ["TOTAL_ASSETS", "TOTAL_PARENT_EQUITY", "ACCOUNTS_RECE", - "OTHER_RECE", "GOODWILL", "SHARE_CAPITAL", "SHORT_LOAN", - "SHORT_FIN_PAYABLE", "NONCURRENT_LIAB_1YEAR", "LONG_LOAN", - "BOND_PAYABLE", "LEASE_LIAB", - # P1 新增存量列 - "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", - "LONG_LOAN", "BOND_PAYABLE", "LEASE_LIAB"] -# DA 组装: cashflow 间接法补充资料四件,列/值缺失按 0(「列存在才加」; -# USERIGHT_ASSET_AMORTIZE 2019 前整列缺失不影响早年 DA) -_DA_PARTS = ["fa_depr", "ia_amort", "lpe_amort", "ua_amort"] -_DATE_COLS = ["REPORT_DATE", "NOTICE_DATE", "UPDATE_DATE"] - -# 输出特征列(P0 32 + P1 35 因子及 2 变体的全部原料;契约由 test_fundamental_library -# 与 test_fundamental_p1_library 锁定) -FEATURE_COLUMNS: list[str] = [ - # 报告期级比率(盈利能力 A / 盈利质量 B / 成长 C / 资本结构 E / 预期事件 F) - "roe_ttm", "roe_deduct_ttm", "roa_ttm", "gp_over_assets", "gross_margin", - "net_margin", "cfo_over_assets", - "tacc", "nonrec_ratio", "impairment_ratio", "invest_income_dep", - "receivables_anomaly", "sales_cash_ratio", "other_rece_ratio", - "rev_q_yoy", "np_q_yoy", "growth_scissors", "gm_delta", "roe_delta", - "asset_growth", "nsi", "ibd_ratio", "goodwill_ratio", - "sue_np", "sue_rev", - "forecast_type_score", "forecast_change_pct", - # 估值/资本行为因子的日频原料(表达式层 ÷ close×share_capital) - "np_ttm", "dnp_ttm", "cfo_ttm", "equity", "share_capital", - "acc_invest_cash_ttm", - # ---- P1 批新增(A 盈利 7)---- - "roe_avg", "roa_pretax", "ebit_over_assets", "ebitda_margin", "roic", - "rd_intensity", "sale_expense_ratio", - # ---- P1 批新增(B 质量 8;B12 为哑变量的连续近似乘积)---- - "inventory_anomaly", "cash_ibd_product", "vsig", "vsig_acc", "vsig_cfo", - "da_intensity", "gm_nm_scissors", "profit_streak", - # ---- P1 批新增(C 成长 8)---- - "np_accel", "rev_accel", "nm_delta", "rev_cagr5", "np_cagr5", - "nwc_growth", "invest_growth", "equity_growth", - # ---- P1 批新增(D 估值 6 的原料;EV = close×share_capital + ev_ex_mv; - # forecast_np_annualized 走 forecast 事件流公告日 asof)---- - "ebit_ttm", "ebitda_ttm", "fcf_ttm", "gp_ttm", "ev_ex_mv", - "forecast_np_annualized", - # ---- P1 批新增(E 资本 3)---- - "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 同剔) -_FIN_NULL_COLS = ["tacc", "nonrec_ratio", "impairment_ratio", - "invest_income_dep", "receivables_anomaly", - "sales_cash_ratio", "other_rece_ratio", - "rev_q_yoy", "np_q_yoy", "growth_scissors", "gm_delta", - "roe_delta", "asset_growth", - "rd_intensity", "sale_expense_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", - # 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") - - -# ==================== 读取层 ==================== - -def _vt_to_file_code(vt_symbol: str) -> str | None: - """``600000.SSE`` → ``600000.SH``;无法映射的(ETF/北交)返回 None.""" - code, _, suffix = vt_symbol.partition(".") - file_suffix = _VT_TO_FILE_SUFFIX.get(suffix.upper()) - return f"{code}.{file_suffix}" if file_suffix else None - - -def _read_static(path: str, want_cols: list[str]) -> pl.DataFrame | None: - """读单股单表 parquet,缺列补 null;零行/缺文件/坏文件 → None(容错跳过).""" - if not os.path.exists(path): - return None - try: - schema = pl.read_parquet_schema(path) - if "REPORT_DATE" not in schema: - return None - cols = [c for c in _DATE_COLS + want_cols if c in schema] - df = pl.read_parquet(path, columns=cols) - except Exception: - return None - if df.height == 0: - return None - missing = [c for c in _DATE_COLS + want_cols if c not in df.columns] - if missing: - df = df.with_columns([pl.lit(None, dtype=pl.Utf8).alias(c) for c in missing]) - for c in want_cols: - df = df.with_columns(pl.col(c).cast(pl.Float64, strict=False)) - return df - - -def _norm_dates(df: pl.DataFrame, cols: list[str]) -> pl.DataFrame: - """日期列归一为 pl.Date:字符串截前 10 位解析,Date/Datetime 直接 cast.""" - exprs = [] - for c in cols: - dtype = df.schema[c] - if dtype == pl.Utf8: - exprs.append(pl.col(c).str.slice(0, 10).str.to_date("%Y-%m-%d", strict=False).alias(c)) - elif dtype != pl.Date: - exprs.append(pl.col(c).cast(pl.Date, strict=False).alias(c)) - return df.with_columns(exprs) if exprs else df - - -def _dedupe_reports(df: pl.DataFrame, value_cols: list[str]) -> pl.DataFrame: - """按 (vt_symbol, REPORT_DATE) 去重:值取 (UPDATE_DATE, NOTICE_DATE) 排序末行 - (重述取终值),有效披露日取组内 NOTICE_DATE 最大值(保守,不提前).""" - df = df.sort(["vt_symbol", "REPORT_DATE", "UPDATE_DATE", "NOTICE_DATE"], - nulls_last=False) - aggs = [pl.col(c).last().alias(c) for c in value_cols] - aggs.append(pl.col("NOTICE_DATE").max().alias("notice_eff")) - return df.group_by(["vt_symbol", "REPORT_DATE"]).agg(aggs) - - -# 各表显式读取清单(income/cashflow 列集不相交,join 不加后缀) -_TABLE_RAW = { - "income": ["TOTAL_OPERATE_INCOME", "OPERATE_COST", "PARENT_NETPROFIT", - "DEDUCT_PARENT_NETPROFIT", "TOTAL_PROFIT", "INVEST_INCOME", - "FAIRVALUE_CHANGE_INCOME", "ASSET_IMPAIRMENT_LOSS", - "CREDIT_IMPAIRMENT_LOSS", - "RESEARCH_EXPENSE", "SALE_EXPENSE", "FE_INTEREST_EXPENSE", - "INCOME_TAX", "BASIC_EPS"], - "balance": _BALANCE_COLS, - "cashflow": ["NETCASH_OPERATE", "SALES_SERVICES", "ACCEPT_INVEST_CASH", - "FA_IR_DEPR", "IA_AMORTIZE", "LPE_AMORTIZE", - "USERIGHT_ASSET_AMORTIZE", "CONSTRUCT_LONG_ASSET", - "RECEIVE_LOAN_CASH", "ISSUE_BOND", "PAY_DEBT_CASH"], -} - - -def _load_statements(codes: list[str], data_dir: str) -> pl.DataFrame: - """全 codes 三表 → 报告期宽表(每 vt_symbol × REPORT_DATE 一行). - - income/cashflow 值列重命名为短名(_CUM_MAP),balance 保持源列名; - anchor = 三表报告期并集,left join 保证单表缺期不拖垮其它表特征。 - """ - table_cols = _TABLE_RAW - per_table: dict[str, pl.DataFrame] = {} - for table, raw_cols in table_cols.items(): - frames = [] - for vt in codes: - file_code = _vt_to_file_code(vt) - if file_code is None: - continue - df = _read_static( - os.path.join(data_dir, table, f"{file_code}_{table}.parquet"), raw_cols) - if df is None: - continue - df = _norm_dates(df, _DATE_COLS) - # 先统一列序/列集再入列(各股文件 schema 子集不同,concat 前必须对齐) - frames.append( - df.with_columns(pl.lit(vt).alias("vt_symbol")) - .select(["vt_symbol", *_DATE_COLS, *raw_cols])) - if not frames: - continue - merged = pl.concat(frames) - if table == "balance": - short = {c: c for c in raw_cols} - else: - short = {k: v for k, v in _CUM_MAP.items() if k in raw_cols} - dedup = _dedupe_reports(merged, list(raw_cols)).rename(short) - per_table[table] = dedup - - empty = pl.DataFrame(schema={"vt_symbol": pl.Utf8, "REPORT_DATE": pl.Date}) - if not per_table: - return empty - - out = pl.concat([t.select(["vt_symbol", "REPORT_DATE"]) - for t in per_table.values()]).unique() - notice_cols = [] - for table, dedup in per_table.items(): - renamed = dedup.rename({"notice_eff": f"_notice_{table}"}) - notice_cols.append(f"_notice_{table}") - out = out.join(renamed, on=["vt_symbol", "REPORT_DATE"], how="left") - # 有效披露日 = 三表 NOTICE_DATE 行最大(最晚可见,保守不提前) - return out.with_columns(pl.max_horizontal(notice_cols).alias("notice_eff")) - - -# ==================== 报告期级指标预计算 ==================== -# 所有 shift/rolling 必须在 .over(_SYM) 组内执行(跨股串行 = 致命错误)。 - -def _qidx() -> pl.Expr: - """连续季度索引: year*4 + quarter(月 3/6/9/12 → 1/2/3/4).""" - return (pl.col("REPORT_DATE").dt.year() * 4 - + (pl.col("REPORT_DATE").dt.month() - 1) // 3 + 1) - - -def _single_quarter(col: str) -> pl.Expr: - """R1 单季化: Q1 直接取累计;其余要求上期恰为上一季度(同年)做差,否则 NaN.""" - cur = pl.col(col) - prev = cur.shift(1) - prev_ok = (pl.col("_qidx") - pl.col("_qidx").shift(1)) == 1 - return ( - pl.when(cur.is_null()).then(None) - .when(pl.col("REPORT_DATE").dt.month() == 3).then(cur) - .when(prev_ok & prev.is_not_null()).then(cur - prev) - .otherwise(None) - ).over(_SYM) - - -def _ttm_of(col: str) -> pl.Expr: - """R2 TTM = 连续 4 个报告期单季之和(窗内任一单季 NaN → NaN).""" - q = pl.col(col) - window_ok = (pl.col("_qidx") - pl.col("_qidx").shift(3)) == 3 - s = q + q.shift(1) + q.shift(2) + q.shift(3) - return pl.when(window_ok).then(s).otherwise(None).over(_SYM) - - -def _yoy4(col: str) -> pl.Expr: - """yoy = X_t / X_{t−4季} − 1;基期缺失/≤0 → NaN(负基数 yoy 无意义).""" - cur, base = pl.col(col), pl.col(col).shift(4) - ok = (pl.col("_qidx") - pl.col("_qidx").shift(4)) == 4 - return ( - pl.when(ok & base.is_not_null() & (base > 0) & cur.is_not_null()) - .then(cur / base - 1.0).otherwise(None) - ).over(_SYM) - - -def _delta4(col: str) -> pl.Expr: - """ΔX = X_t − X_{t−4季}(要求恰好隔 4 个季度).""" - ok = (pl.col("_qidx") - pl.col("_qidx").shift(4)) == 4 - return pl.when(ok).then(pl.col(col) - pl.col(col).shift(4)).otherwise(None).over(_SYM) - - -def _mean4(col: str) -> pl.Expr: - """mean(X_t, X_{t−4季})(A03 平均 ROE 分母;恰隔 4 季 + 两期非空守卫).""" - ok = (pl.col("_qidx") - pl.col("_qidx").shift(4)) == 4 - both = pl.col(col).is_not_null() & pl.col(col).shift(4).is_not_null() - m = (pl.col(col) + pl.col(col).shift(4)) / 2.0 - return pl.when(ok & both).then(m).otherwise(None).over(_SYM) - - -def _cagr5(col: str) -> pl.Expr: - """5 年 CAGR = (X_y / X_{y−5})^{1/5} − 1(年报行;基期/现期 ≤0 → NaN, - 负基期 CAGR 无意义;恰隔 20 季守卫).""" - cur, base = pl.col(col), pl.col(col).shift(20) - ok = (pl.col("_qidx") - pl.col("_qidx").shift(20)) == 20 - valid = ok & base.is_not_null() & (base > 0) & cur.is_not_null() & (cur > 0) - return pl.when(valid).then((cur / base).pow(0.2) - 1.0).otherwise(None).over(_SYM) - - -def _std16(col: str) -> pl.Expr: - """16 季滚动 sample std(ddof=1,与 SUE Foster 同款钉死);窗口须恰为连续 - 16 个季度且全非空——不足 16 期/窗内含缺失 → NaN(不填 0).""" - ok = ((pl.col("_qidx") - pl.col("_qidx").shift(15)) == 15).over(_SYM) - sd = pl.col(col).rolling_std(window_size=16, ddof=1).over(_SYM) - return pl.when(ok & sd.is_not_null()).then(sd).otherwise(None) - - -def _safe_ratio(num: pl.Expr, den: pl.Expr) -> pl.Expr: - """分母缺失/为 0 → NaN(比率类通用守卫).""" - return pl.when(den.is_not_null() & (den != 0)).then(num / den).otherwise(None) - - -def _compute_report_features(reports: pl.DataFrame) -> pl.DataFrame: - """报告期宽表 → 全部报告期级特征(逐级 with_columns,over 组内时序).""" - df = reports.sort([_SYM, "REPORT_DATE"]).with_columns(_qidx().alias("_qidx")) - - # 第一级: 单季化 + TTM(累计列) - df = df.with_columns( - [_single_quarter(c).alias(f"q_{c}") for c in _CUM_MAP.values()] - ).with_columns( - [_ttm_of(f"q_{c}").alias(f"ttm_{c}") for c in _CUM_MAP.values()] - ) - - # 第一级半(P1): 补充资料加工底座 - # - DA 累计 = 四件折旧摊销之和(缺列/缺值按 0,与 IBD 组装同款; - # USERIGHT_ASSET_AMORTIZE 2019 前整列缺失不影响早年 DA) - # - 利息费用累计 fill 0(NAS 实测披露稀疏: 600519 78% null/银行模板整列缺; - # 未披露按 0 回加 → EBIT 退化为 TP,口径登记) - # - NWC = 存货+应收(null 传播,缺一即 NaN) - # - 年报行门控列(C13/C14/C18 年度口径;非年报行 null → yoy 天然 NaN) - df = df.with_columns( - pl.sum_horizontal([pl.col(p).fill_null(0.0) for p in _DA_PARTS]).alias("_da_cum"), - pl.col("int_exp").fill_null(0.0).alias("_int0_cum"), - (pl.col("INVENTORY") + pl.col("ACCOUNTS_RECE")).alias("_nwc"), - pl.when(pl.col("REPORT_DATE").dt.month() == 12).then(pl.col("rev")).alias("_rev_ann"), - pl.when(pl.col("REPORT_DATE").dt.month() == 12).then(pl.col("np")).alias("_np_ann"), - pl.when(pl.col("REPORT_DATE").dt.month() == 12).then(pl.col("capex")).alias("_capex_ann"), - ).with_columns( - _single_quarter("_da_cum").alias("q__da"), - _single_quarter("_int0_cum").alias("q__int0"), - ).with_columns( - _ttm_of("q__da").alias("ttm__da"), - _ttm_of("q__int0").alias("ttm__int0"), - ) - - # 第二级: IBD(缺组件按 0,钉死含 LEASE_LIAB 版) + 金融股判定(银行模板无 - # 营业成本) + 毛利 + τ 实际税率(TTM 口径: 消费者均为 TTM 流量; - # τ = INCOME_TAX_TTM/TOTAL_PROFIT_TTM 截断 [0,0.5],两列缺失或 TP≤0 → 0.25) - ibd = pl.sum_horizontal([pl.col(p).fill_null(0.0) for p in _IBD_PARTS]) - is_fin = pl.col("cogs").is_null() | (pl.col("cogs") == 0) - tau = ( - pl.when(pl.col("ttm_tp").is_not_null() & (pl.col("ttm_tp") > 0) - & pl.col("ttm_tax").is_not_null()) - .then((pl.col("ttm_tax") / pl.col("ttm_tp")).clip(0.0, 0.5)) - .otherwise(0.25) - ) - df = df.with_columns( - ibd.alias("_ibd"), - is_fin.alias("_is_fin"), - tau.alias("_tau"), - (pl.col("ttm_rev") - pl.col("ttm_cogs")).alias("_gp_ttm"), - # EBIT_TTM = (TOTAL_PROFIT + FE_INTEREST_EXPENSE)_TTM(§1.3) - (pl.col("ttm_tp") + pl.col("ttm__int0")).alias("_ebit_ttm"), - # (EBIT+DA)_TTM / FCF_TTM / 债务净发行 TTM(E03 三流合成) - (pl.col("ttm_tp") + pl.col("ttm__int0") + pl.col("ttm__da")).alias("_ebitda_ttm"), - (pl.col("ttm_cfo") - pl.col("ttm_capex")).alias("_fcf_ttm"), - (pl.col("ttm_recv_loan") + pl.col("ttm_issue_bond") - - pl.col("ttm_pay_debt")).alias("_debt_issue_ttm"), - ) - - # 第二级半(P1): 平均净资产(A03 分母) + VSIG 三序列底座(单季口径 / TA) - ta = pl.col("TOTAL_ASSETS") - df = df.with_columns( - _mean4("TOTAL_PARENT_EQUITY").alias("_eq_avg"), - _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) - eq = pl.col("TOTAL_PARENT_EQUITY") - df = df.with_columns( - # 盈利能力 A - _safe_ratio(pl.col("ttm_np"), eq).alias("roe_ttm"), - _safe_ratio(pl.col("ttm_dnp"), eq).alias("roe_deduct_ttm"), - _safe_ratio(pl.col("ttm_np"), ta).alias("roa_ttm"), - _safe_ratio(pl.col("_gp_ttm"), ta).alias("gp_over_assets"), - _safe_ratio(pl.col("_gp_ttm"), pl.col("ttm_rev")).alias("gross_margin"), - _safe_ratio(pl.col("ttm_np"), pl.col("ttm_rev")).alias("net_margin"), - _safe_ratio(pl.col("ttm_cfo"), ta).alias("cfo_over_assets"), - # 盈利质量 B - _safe_ratio(pl.col("ttm_np") - pl.col("ttm_cfo"), ta).alias("tacc"), - _safe_ratio(pl.col("ttm_np") - pl.col("ttm_dnp"), - pl.col("ttm_np").abs()).alias("nonrec_ratio"), - _safe_ratio((pl.col("ttm_asset_imp") + pl.col("ttm_credit_imp")).abs(), - ta).alias("impairment_ratio"), - _safe_ratio(pl.col("ttm_invest_inc") + pl.col("ttm_fv_inc"), - pl.col("ttm_tp").abs()).alias("invest_income_dep"), - _safe_ratio(pl.col("ttm_sales_cash"), pl.col("ttm_rev")).alias("sales_cash_ratio"), - _safe_ratio(pl.col("OTHER_RECE"), ta).alias("other_rece_ratio"), - # 成长 C / 资本结构 E - _safe_ratio(pl.col("_ibd"), ta).alias("ibd_ratio"), - _safe_ratio(pl.col("GOODWILL"), ta).alias("goodwill_ratio"), - # 盈利能力 A(P1 7): A03/A09/A10/A11/A12/A15/A16 - _safe_ratio(pl.col("ttm_np"), pl.col("_eq_avg")).alias("roe_avg"), - _safe_ratio(pl.col("ttm_np") + pl.col("ttm__int0") * (1.0 - pl.col("_tau")), - ta).alias("roa_pretax"), - _safe_ratio(pl.col("_ebit_ttm"), ta).alias("ebit_over_assets"), - _safe_ratio(pl.col("_ebitda_ttm"), pl.col("ttm_rev")).alias("ebitda_margin"), - _safe_ratio(pl.col("_ebit_ttm") * (1.0 - pl.col("_tau")), - eq + pl.col("_ibd") - pl.col("MONETARYFUNDS")).alias("roic"), - _safe_ratio(pl.col("ttm_research"), pl.col("ttm_rev")).alias("rd_intensity"), - _safe_ratio(pl.col("ttm_sale_exp"), pl.col("ttm_rev")).alias("sale_expense_ratio"), - # 盈利质量 B(P1): B12 哑变量的连续近似 = (MON/TA)×(IBD/TA) 乘积变体 - # (表达式引擎无截面分位函数,不改引擎——survey B12 的可计算降级) - (_safe_ratio(pl.col("MONETARYFUNDS"), ta) - * _safe_ratio(pl.col("_ibd"), ta)).alias("cash_ibd_product"), - _safe_ratio(pl.col("ttm__da"), pl.col("ttm_rev")).alias("da_intensity"), - # E07 利息保障倍数(利息费用≤0 → NaN: 负利息=净收入,倍数无意义) - 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 组内时序) - df = df.with_columns( - _yoy4("q_rev").alias("rev_q_yoy"), - _yoy4("q_np").alias("np_q_yoy"), - _yoy4("ACCOUNTS_RECE").alias("_ar_yoy"), - _yoy4("ttm_rev").alias("_rev_ttm_yoy"), - _yoy4("SHARE_CAPITAL").alias("nsi"), - _yoy4("TOTAL_ASSETS").alias("asset_growth"), - _delta4("gross_margin").alias("gm_delta"), - _delta4("roe_ttm").alias("roe_delta"), - _delta4("q_np").alias("_diff4_np"), - _delta4("q_rev").alias("_diff4_rev"), - # P1: B09 存货同比 / C16 NWC 同比 / C19 净资产 / E10 商誉 - _yoy4("INVENTORY").alias("_inv_yoy"), - _yoy4("_nwc").alias("nwc_growth"), - _yoy4("TOTAL_PARENT_EQUITY").alias("equity_growth"), - _yoy4("GOODWILL").alias("goodwill_growth"), - # P1: C18 投资增速(年度口径,非年报行 cur=null → NaN) - _yoy4("_capex_ann").alias("invest_growth"), - # P1: C13/C14 五年 CAGR(年报行,恰隔 20 季守卫,基期/现期≤0 → NaN) - _cagr5("_rev_ann").alias("rev_cagr5"), - _cagr5("_np_ann").alias("np_cagr5"), - _delta4("net_margin").alias("nm_delta"), - _delta4("q_eps").alias("_diff4_eps"), - ).with_columns( - (pl.col("np_q_yoy") - pl.col("rev_q_yoy")).alias("growth_scissors"), - (pl.col("_ar_yoy") - pl.col("_rev_ttm_yoy")).alias("receivables_anomaly"), - # B09 与 B08 同构: 期末存量同比 − REV_TTM 同比(登记口径) - (pl.col("_inv_yoy") - pl.col("_rev_ttm_yoy")).alias("inventory_anomaly"), - # B18 毛净剪刀差 = GM_TTM − NM_TTM(第三级产物,同块不可引用故后置) - (pl.col("gross_margin") - pl.col("net_margin")).alias("gm_nm_scissors"), - # C07/C08 加速度 = yoy 的恰隔 4 季二次差分(基期>0 守卫由 yoy 层继承, - # 二次差分同样 NaN 传播;同块不可引用 yoy 列故后置) - _delta4("np_q_yoy").alias("np_accel"), - _delta4("rev_q_yoy").alias("rev_accel"), - ) - - # 第五级: SUE(Foster 标准化)= diff4 / std(过去 8 期 diff4, ddof=1) - for src, out in (("_diff4_np", "sue_np"), ("_diff4_rev", "sue_rev"), - ("_diff4_eps", "sue_eps")): - sd = pl.col(src).rolling_std(window_size=8, ddof=1).over(_SYM) - df = df.with_columns( - pl.when(sd.is_not_null() & (sd > 0) & pl.col(src).is_not_null()) - .then(pl.col(src) / sd).otherwise(None).alias(out) - ) - # SUE 严窗变体(随批互评): σ 只用 t−1 及更早差分(shift(1) 后滚 8 期,不含当期) - sd_strict = pl.col("_diff4_np").shift(1).rolling_std(window_size=8, ddof=1).over(_SYM) - df = df.with_columns( - pl.when(sd_strict.is_not_null() & (sd_strict > 0) & pl.col("_diff4_np").is_not_null()) - .then(pl.col("_diff4_np") / sd_strict).otherwise(None).alias("sue_np_strict") - ) - - # 第五级半(P1): VSIG 16 季滚动(sample std ddof=1 钉死,恰连续 16 季全非空) - # + F06 披露及时性 = −(有效披露日 − 报告期) 天数(早披露=高分;notice_eff - # 取三表最晚可见,与 PIT 锚一致) - df = df.with_columns( - _std16("_np_ta").alias("vsig"), - _std16("_accq_ta").alias("vsig_acc"), - _std16("_cfo_ta").alias("vsig_cfo"), - (-(pl.col("notice_eff") - pl.col("REPORT_DATE")).dt.total_days()) - .cast(pl.Float64).alias("disclosure_speed"), - ) - - # 第五级半续(P1): B19 持续盈利季数 = 连续单季 NP>0 计数(截断 8; - # 当期缺失→NaN;非正→0 断流;中间缺失行视为断流点) - df = df.with_columns( - (pl.col("q_np") > 0).alias("_pos"), - pl.int_range(pl.len()).cast(pl.Int64).alias("_ridx"), - ).with_columns( - pl.when(pl.col("_pos").is_null() | ~pl.col("_pos")) - .then(pl.col("_ridx")).otherwise(None) - .fill_null(strategy="forward").over(_SYM).alias("_lastbrk"), - ).with_columns( - pl.when(pl.col("_pos").is_null()).then(None) - .when(pl.col("_pos")) - .then((pl.col("_ridx") - pl.col("_lastbrk").fill_null(-1)).clip(1, 8).cast(pl.Float64)) - .otherwise(0.0).alias("profit_streak"), - ) - - # 金融股: 盈利质量/成长/费用类族特征置 NaN(§7 红线 5;估值族保留) - df = df.with_columns([ - pl.when(pl.col("_is_fin")).then(None).otherwise(pl.col(c)).alias(c) - for c in _FIN_NULL_COLS - ]) - - # 输出别名(估值/资本行为因子的日频原料) - df = df.with_columns( - pl.col("ttm_np").alias("np_ttm"), - pl.col("ttm_dnp").alias("dnp_ttm"), - pl.col("ttm_cfo").alias("cfo_ttm"), - pl.col("ttm_acc_inv_cash").alias("acc_invest_cash_ttm"), - pl.col("TOTAL_PARENT_EQUITY").alias("equity"), - pl.col("SHARE_CAPITAL").alias("share_capital"), - # P1: EV 群/FCF/债务净发行原料(EV = close×share_capital + ev_ex_mv) - pl.col("_gp_ttm").alias("gp_ttm"), - pl.col("_ebit_ttm").alias("ebit_ttm"), - pl.col("_ebitda_ttm").alias("ebitda_ttm"), - pl.col("_fcf_ttm").alias("fcf_ttm"), - (pl.col("_ibd") - pl.col("MONETARYFUNDS")).alias("ev_ex_mv"), - pl.col("_debt_issue_ttm").alias("debt_issue_ttm"), - ) - return df - - -# ==================== forecast 事件层 ==================== - -# 预告净利年化系数(按报告期进度;D13 预期 EP): Q1×4 / H1×2 / Q3×4/3 / 年报×1 -_ANNUALIZE_FACTOR = {3: 4.0, 6: 2.0, 9: 4.0 / 3.0, 12: 1.0} - - -def _load_forecast_events(codes: list[str], data_dir: str) -> tuple[pl.DataFrame, pl.DataFrame]: - """forecast 按期文件 → (fc_events, fc_pair). - - fc_events: (vt_symbol, eff=公告日期, forecast_type_score, forecast_change_pct, - forecast_np_annualized) 事件行——一股一公告日多行(按预测指标), - 归母净利润行优先(含"净利润"且不含"扣"),无净利润行 fallback 任意行; - 年化预告净利只对净利行生效(fallback 营业收入行的中值不作净利用)。 - 同股多公告日全保留(asof 取最新)。 - - fc_pair: (vt_symbol, REPORT_DATE, _fc_mid, eff) —— F07 预告兑现差的配对原料, - 每 (股, 报告期) 取最新公告日的净利行中值(REPORT_DATE 取自文件名)。 - """ - schema = {"vt_symbol": pl.Utf8, "eff": pl.Date, "REPORT_DATE": pl.Date, - "forecast_type_score": pl.Float64, "forecast_change_pct": pl.Float64, - "forecast_np_annualized": pl.Float64, "_fc_mid": pl.Float64, "_is_np": pl.Boolean} - fc_dir = os.path.join(data_dir, "forecast") - if not os.path.isdir(fc_dir): - return pl.DataFrame(schema=schema), pl.DataFrame(schema=schema) - code_set = set(codes) - frames = [] - for fname in sorted(os.listdir(fc_dir)): - if not fname.endswith(".parquet"): - continue - try: - f = pl.read_parquet(os.path.join(fc_dir, fname)) - except Exception: - continue - if f.height == 0 or not all(c in f.columns for c in - ("股票代码", "预告类型", "公告日期")): - continue - try: # 文件名前 8 位 = 报告期(20230630_forecast.parquet) - report_date = datetime.strptime(fname[:8], "%Y%m%d").date() - except ValueError: - continue - code = pl.col("股票代码").cast(pl.Utf8).str.strip_chars().str.zfill(6) - # 交易所映射: 60→SSE;北交前缀白名单(92/43/82/83)→BJSE;其余→SZSE - # (互评备注: 北交种类不得落入 SZSE——容器/实盘 universe 按后缀路由) - is_bj = (code.str.starts_with("92") | code.str.starts_with("43") - | code.str.starts_with("82") | code.str.starts_with("83")) - vt = (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")).alias("vt_symbol")) - if "预测指标" in f.columns: - ind = pl.col("预测指标").cast(pl.Utf8) - pref = (ind.str.contains("净利润") & ~ind.str.contains("扣")).cast(pl.Int32) - else: - pref = pl.lit(0, pl.Int32) - pct = (pl.col("业绩变动幅度").cast(pl.Float64, strict=False) - if "业绩变动幅度" in f.columns else pl.lit(None, pl.Float64)) - mid = (pl.col("预测数值").cast(pl.Float64, strict=False) - if "预测数值" in f.columns else pl.lit(None, pl.Float64)) - f = f.with_columns( - vt, - pref.alias("_pref"), - pct.alias("_pct"), - mid.alias("_mid"), - pl.lit(report_date, dtype=pl.Date).alias("REPORT_DATE"), - pl.lit(_ANNUALIZE_FACTOR.get(report_date.month), dtype=pl.Float64).alias("_annf"), - pl.col("公告日期").cast(pl.Date, strict=False).alias("eff"), - pl.col("预告类型").cast(pl.Utf8).replace_strict( - FORECAST_TYPE_SCORE, default=None, return_dtype=pl.Float64 - ).alias("_score"), - ).filter(pl.col("vt_symbol").is_in(code_set) & pl.col("eff").is_not_null()) - if f.height: - frames.append(f.select( - ["vt_symbol", "eff", "REPORT_DATE", "_pref", "_score", "_pct", "_mid", "_annf"])) - if not frames: - return pl.DataFrame(schema=schema), pl.DataFrame(schema=schema) - fc = pl.concat(frames).sort(["vt_symbol", "eff", "_pref"]) - # 同 (vt, 公告日, 报告期) 取优先级最高行(_pref 大者排序在后 → last); - # 年化预告净利 = 净利行中值 × 年化系数(非净利行 fallback → null) - events = fc.group_by(["vt_symbol", "eff", "REPORT_DATE"]).agg( - pl.col("_score").last().alias("forecast_type_score"), - pl.col("_pct").last().alias("forecast_change_pct"), - pl.col("_pref").last().alias("_is_np"), - pl.col("_mid").last().alias("_mid"), - pl.col("_annf").last().alias("_annf"), - ).with_columns( - pl.when(pl.col("_is_np") == 1) - .then(pl.col("_mid") * pl.col("_annf")).otherwise(None) - .alias("forecast_np_annualized"), - ) - # F07 配对: 每 (股, 报告期) 最新公告日的净利行中值 - fc_pair = (events.filter(pl.col("_is_np") == 1 & pl.col("_mid").is_not_null()) - .sort(["vt_symbol", "REPORT_DATE", "eff"]) - .group_by(["vt_symbol", "REPORT_DATE"]).agg( - pl.col("eff").last().alias("eff"), - pl.col("_mid").last().alias("_fc_mid")) - .select(["vt_symbol", "REPORT_DATE", "_fc_mid", "eff"])) - # 事件流: 同 (股, 公告日) 多报告期行罕见(同年同日两期预告)——取最新报告期 - # 为当前信号(P0 语义 = 每公告日一行) - fc_events = (events.sort(["vt_symbol", "eff", "REPORT_DATE"]) - .group_by(["vt_symbol", "eff"]).last() - .select(["vt_symbol", "eff", "forecast_type_score", - "forecast_change_pct", "forecast_np_annualized"])) - 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)单次 -# join_asof + 全量副本必 OOM;按股分批使峰值 ≈ 事件表(常驻) + 单批 grid + 输出累积 -BATCH_CODES = 500 - - -def _prepare_days(start: str, end: str, trading_dates) -> list: - """交易日/日历日序列(一次解析,各批共用).""" - if trading_dates is not None: - days = pl.Series("datetime", trading_dates).cast(pl.Datetime("us")) - return days.unique().sort().to_list() - s = datetime.strptime(start, "%Y-%m-%d") - e = datetime.strptime(end, "%Y-%m-%d") - return pl.datetime_range(s, e, interval="1d", eager=True).cast( - pl.Datetime("us")).to_list() - - -def _build_grid(codes: list[str], day_list: list) -> pl.DataFrame: - """单批日频 grid: codes × day_list(批内全量,批间无全量 grid 副本).""" - return pl.DataFrame({ - "vt_symbol": pl.Series([c for c in codes for _ in day_list], dtype=pl.Utf8), - "datetime": pl.Series(day_list * len(codes), dtype=pl.Datetime("us")), - }, schema={"vt_symbol": pl.Utf8, "datetime": pl.Datetime("us")}) - - -def _load_feature_events(codes: list[str], data_dir: str) -> tuple[pl.DataFrame, pl.DataFrame, 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 - and c not in _DOMAIN_COLS] - reports = _load_statements(codes, data_dir) - if reports.height == 0: - stmt_events = pl.DataFrame(schema={ - "vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), **{c: pl.Float64 for c in stmt_cols}}) - beat_events = pl.DataFrame(schema={ - "vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "forecast_beat": pl.Float64}) - else: - feat = _compute_report_features(reports) - # NOTICE_DATE 缺失报告期整期跳过(红线: 宁缺毋假) - feat = feat.filter(pl.col("notice_eff").is_not_null()) - stmt_events = feat.select( - pl.col("vt_symbol"), - pl.col("notice_eff").cast(pl.Datetime("us")).alias("eff"), - *stmt_cols, - ).sort("eff") - beat_events = _build_beat_events(feat, codes, data_dir) - fc_events, _ = _load_forecast_events(codes, data_dir) - fc_events = fc_events.select( - pl.col("vt_symbol"), - pl.col("eff").cast(pl.Datetime("us")), - *_FORECAST_COLS, - ).sort("eff") - return stmt_events, fc_events, beat_events - - -def _build_beat_events(feat: pl.DataFrame, codes: list[str], data_dir: str) -> pl.DataFrame: - """F07 预告兑现差事件流: (实际NP − 预告中值)/abs(预告中值),同 REPORT_DATE 配对. - - 前视红线(P1 任务书): 兑现差含实际 NP,只有实际报告披露后才可知—— - PIT 锚 = max(该报告期 income 有效披露日 notice_eff, 预告公告日), - 不早于两者较晚者(预告公告晚于年报的罕见情形不被提前泄露)。 - 实际 NP 用报告期累计归母净利(预告口径即期间累计);预告中值=0/缺 → NaN。 - """ - schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "forecast_beat": pl.Float64} - _, fc_pair = _load_forecast_events(codes, data_dir) - if fc_pair.height == 0: - return pl.DataFrame(schema=schema) - beat = ( - feat.select("vt_symbol", "REPORT_DATE", "notice_eff", pl.col("np")) - .join(fc_pair, on=["vt_symbol", "REPORT_DATE"], how="inner") - .with_columns( - _safe_ratio(pl.col("np") - pl.col("_fc_mid"), - pl.col("_fc_mid").abs()).alias("forecast_beat"), - pl.max_horizontal("notice_eff", "eff").alias("_anchor"), - ) - .filter(pl.col("forecast_beat").is_not_null() & pl.col("_anchor").is_not_null()) - ) - return beat.select( - pl.col("vt_symbol"), - pl.col("_anchor").cast(pl.Datetime("us")).alias("eff"), - pl.col("forecast_beat"), - ).sort("eff") - - -def iter_fundamental_feature_chunks( - codes: list[str], - start: str, - end: str, - data_dir: str = DEFAULT_STATIC_DIR, - trading_dates: pl.Series | list | None = None, - batch_codes: int = BATCH_CODES, - columns: list[str] | None = None, -): - """按 vt_symbol 分批产出 PIT 日频特征块(生成器,NAS 全量防 OOM 主入口). - - 每块 = 一批 codes × 全部日期 × columns(默认 FEATURE_COLUMNS 全量), - 顺序即 codes 列表顺序;批内 grid 用完即弃,事件右表(报告期+forecast+ - 兑现差)全批共用仅此一份。batch_eval 侧应逐块 join alpha_df 分片后 - concat,避免持有本帧全量副本;columns 子集可只 join 本批表达式引用列, - 全量 68 列 × 1480 万行 ≈ 8G——按引用瘦身是 NAS 7.9G 内存的关键杠杆。 - """ - 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] - grid = _build_grid(chunk_codes, day_list).sort("datetime") - # join_asof 前已显式按键排序;polars 1.42 用 by 分组时无法校验 sortedness, - # 该提示无信息量,就地抑制(sort 即正确性保险) - with warnings.catch_warnings(): - warnings.simplefilter("ignore", UserWarning) - # 三条事件流各自 asof 后丢弃右表键 eff(留置会以 eff_right 后缀 - # 累积,第三次 join 撞名) - out = grid.join_asof( - stmt_events, left_on="datetime", right_on="eff", - by="vt_symbol", strategy="backward").drop("eff") - out = out.sort("datetime").join_asof( - fc_events, left_on="datetime", right_on="eff", - by="vt_symbol", strategy="backward").drop("eff") - 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 # 批间释放(下一批重绑定) - - -def build_fundamental_features( - codes: list[str], - start: str, - end: str, - data_dir: str = DEFAULT_STATIC_DIR, - trading_dates: pl.Series | list | None = None, - batch_codes: int = BATCH_CODES, - columns: list[str] | None = None, -) -> pl.DataFrame: - """构建 PIT 日频财务特征: vt_symbol × datetime × columns. - - Args: - codes: vt_symbol 列表(如 "600000.SSE") - start/end: "YYYY-MM-DD" 窗口(trading_dates=None 时生成日历日 grid) - data_dir: 静态域根目录(NAS=/volume1/stock/sanguo_vnpy_v2/data/static) - trading_dates: 交易日子集(传 bars 的 unique datetime 免造非交易日行) - batch_codes: 按股分批大小(全量防 OOM;测试可调小验分块等值) - columns: 输出特征列子集(默认 FEATURE_COLUMNS 全量;引用瘦身用) - - Returns: - 每行 = 决策日可见的最新报告期特征(NOTICE_DATE ≤ 决策日,asof 前向填充)。 - """ - out_cols = list(columns) if columns is not None else list(FEATURE_COLUMNS) - schema = {"vt_symbol": pl.Utf8, "datetime": pl.Datetime("us"), - **{c: pl.Float64 for c in out_cols}} - if not codes: - return pl.DataFrame(schema=schema) - return pl.concat( - iter_fundamental_feature_chunks( - codes, start, end, data_dir=data_dir, - trading_dates=trading_dates, batch_codes=batch_codes, columns=out_cols), - how="vertical") - +from .fundamental_schema import ( # noqa: F401 + FEATURE_COLUMNS, + FORECAST_TYPE_SCORE, + DEFAULT_STATIC_DIR, + _BALANCE_COLS, + _BEAT_COLS, + _CUM_MAP, + _DA_PARTS, + _DATE_COLS, + _DOMAIN_COLS, + _FIN_NULL_COLS, + _FORECAST_COLS, + _IBD_PARTS, + _SYM, + _VT_TO_FILE_SUFFIX, +) +from .fundamental_statements import ( # noqa: F401 + _TABLE_RAW, + _dedupe_reports, + _load_statements, + _norm_dates, + _read_static, + _vt_to_file_code, +) +from .fundamental_report_features import ( # noqa: F401 + _cagr5, + _compute_report_features, + _delta4, + _mean4, + _qidx, + _safe_ratio, + _single_quarter, + _std16, + _ttm_of, + _yoy4, +) +from .fundamental_forecast import ( # noqa: F401 + _ANNUALIZE_FACTOR, + _load_forecast_events, +) +from .fundamental_domains import ( # noqa: F401 + _DIVIDEND_CHAIN, + _EMPTY_EVENTS, + _GDHS_CUTOFF_CANDIDATES, + _VB_GAP, + _code6_to_vt, + _inv_guard, + _load_dividend_cum, + _load_extra_domains, + _load_gdhs_events, + _load_topholder_events, + _load_vb_daily, + _missing_warned, + _warn_domain_missing, +) +from .fundamental_pit import ( # noqa: F401 + BATCH_CODES, + _build_beat_events, + _build_grid, + _load_feature_events, + _prepare_days, + build_fundamental_features, + iter_fundamental_feature_chunks, +) diff --git a/sanguo_factor/fundamental_domains.py b/sanguo_factor/fundamental_domains.py new file mode 100644 index 0000000..8664296 --- /dev/null +++ b/sanguo_factor/fundamental_domains.py @@ -0,0 +1,329 @@ +# sanguo_factor/fundamental_domains.py +"""P1-B 四新域事件层: dividend/gdhs/top_holders/valuation_baostock +(自 fundamental_adapter.py 拆出,纯结构重构). + +四域缺任一 → 相关特征列全 NaN + 一次性 warning(优雅降级,本地可跑通); +_load_extra_domains 按引用列瘦身加载(未引用的域不读盘). +""" +from __future__ import annotations + +import os +import warnings +from datetime import date, datetime + +import polars as pl + +from .fundamental_schema import _SYM +from .fundamental_statements import _norm_dates, _vt_to_file_code + +# ==================== 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), + } diff --git a/sanguo_factor/fundamental_forecast.py b/sanguo_factor/fundamental_forecast.py new file mode 100644 index 0000000..a7e9539 --- /dev/null +++ b/sanguo_factor/fundamental_forecast.py @@ -0,0 +1,118 @@ +# sanguo_factor/fundamental_forecast.py +"""forecast 业绩预告事件层: static/forecast 按报告期全市场文件 → 事件流 +(自 fundamental_adapter.py 拆出,纯结构重构). + +fc_events(公告日 asof 信号)+ fc_pair(F07 预告兑现差配对原料)双出口; +归母净利润行优先,年化预告净利只对净利行生效. +""" +from __future__ import annotations + +import os +from datetime import datetime + +import polars as pl + +from .fundamental_schema import FORECAST_TYPE_SCORE + +# ==================== forecast 事件层 ==================== + +# 预告净利年化系数(按报告期进度;D13 预期 EP): Q1×4 / H1×2 / Q3×4/3 / 年报×1 +_ANNUALIZE_FACTOR = {3: 4.0, 6: 2.0, 9: 4.0 / 3.0, 12: 1.0} + + +def _load_forecast_events(codes: list[str], data_dir: str) -> tuple[pl.DataFrame, pl.DataFrame]: + """forecast 按期文件 → (fc_events, fc_pair). + + fc_events: (vt_symbol, eff=公告日期, forecast_type_score, forecast_change_pct, + forecast_np_annualized) 事件行——一股一公告日多行(按预测指标), + 归母净利润行优先(含"净利润"且不含"扣"),无净利润行 fallback 任意行; + 年化预告净利只对净利行生效(fallback 营业收入行的中值不作净利用)。 + 同股多公告日全保留(asof 取最新)。 + + fc_pair: (vt_symbol, REPORT_DATE, _fc_mid, eff) —— F07 预告兑现差的配对原料, + 每 (股, 报告期) 取最新公告日的净利行中值(REPORT_DATE 取自文件名)。 + """ + schema = {"vt_symbol": pl.Utf8, "eff": pl.Date, "REPORT_DATE": pl.Date, + "forecast_type_score": pl.Float64, "forecast_change_pct": pl.Float64, + "forecast_np_annualized": pl.Float64, "_fc_mid": pl.Float64, "_is_np": pl.Boolean} + fc_dir = os.path.join(data_dir, "forecast") + if not os.path.isdir(fc_dir): + return pl.DataFrame(schema=schema), pl.DataFrame(schema=schema) + code_set = set(codes) + frames = [] + for fname in sorted(os.listdir(fc_dir)): + if not fname.endswith(".parquet"): + continue + try: + f = pl.read_parquet(os.path.join(fc_dir, fname)) + except Exception: + continue + if f.height == 0 or not all(c in f.columns for c in + ("股票代码", "预告类型", "公告日期")): + continue + try: # 文件名前 8 位 = 报告期(20230630_forecast.parquet) + report_date = datetime.strptime(fname[:8], "%Y%m%d").date() + except ValueError: + continue + code = pl.col("股票代码").cast(pl.Utf8).str.strip_chars().str.zfill(6) + # 交易所映射: 60→SSE;北交前缀白名单(92/43/82/83)→BJSE;其余→SZSE + # (互评备注: 北交种类不得落入 SZSE——容器/实盘 universe 按后缀路由) + is_bj = (code.str.starts_with("92") | code.str.starts_with("43") + | code.str.starts_with("82") | code.str.starts_with("83")) + vt = (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")).alias("vt_symbol")) + if "预测指标" in f.columns: + ind = pl.col("预测指标").cast(pl.Utf8) + pref = (ind.str.contains("净利润") & ~ind.str.contains("扣")).cast(pl.Int32) + else: + pref = pl.lit(0, pl.Int32) + pct = (pl.col("业绩变动幅度").cast(pl.Float64, strict=False) + if "业绩变动幅度" in f.columns else pl.lit(None, pl.Float64)) + mid = (pl.col("预测数值").cast(pl.Float64, strict=False) + if "预测数值" in f.columns else pl.lit(None, pl.Float64)) + f = f.with_columns( + vt, + pref.alias("_pref"), + pct.alias("_pct"), + mid.alias("_mid"), + pl.lit(report_date, dtype=pl.Date).alias("REPORT_DATE"), + pl.lit(_ANNUALIZE_FACTOR.get(report_date.month), dtype=pl.Float64).alias("_annf"), + pl.col("公告日期").cast(pl.Date, strict=False).alias("eff"), + pl.col("预告类型").cast(pl.Utf8).replace_strict( + FORECAST_TYPE_SCORE, default=None, return_dtype=pl.Float64 + ).alias("_score"), + ).filter(pl.col("vt_symbol").is_in(code_set) & pl.col("eff").is_not_null()) + if f.height: + frames.append(f.select( + ["vt_symbol", "eff", "REPORT_DATE", "_pref", "_score", "_pct", "_mid", "_annf"])) + if not frames: + return pl.DataFrame(schema=schema), pl.DataFrame(schema=schema) + fc = pl.concat(frames).sort(["vt_symbol", "eff", "_pref"]) + # 同 (vt, 公告日, 报告期) 取优先级最高行(_pref 大者排序在后 → last); + # 年化预告净利 = 净利行中值 × 年化系数(非净利行 fallback → null) + events = fc.group_by(["vt_symbol", "eff", "REPORT_DATE"]).agg( + pl.col("_score").last().alias("forecast_type_score"), + pl.col("_pct").last().alias("forecast_change_pct"), + pl.col("_pref").last().alias("_is_np"), + pl.col("_mid").last().alias("_mid"), + pl.col("_annf").last().alias("_annf"), + ).with_columns( + pl.when(pl.col("_is_np") == 1) + .then(pl.col("_mid") * pl.col("_annf")).otherwise(None) + .alias("forecast_np_annualized"), + ) + # F07 配对: 每 (股, 报告期) 最新公告日的净利行中值 + fc_pair = (events.filter(pl.col("_is_np") == 1 & pl.col("_mid").is_not_null()) + .sort(["vt_symbol", "REPORT_DATE", "eff"]) + .group_by(["vt_symbol", "REPORT_DATE"]).agg( + pl.col("eff").last().alias("eff"), + pl.col("_mid").last().alias("_fc_mid")) + .select(["vt_symbol", "REPORT_DATE", "_fc_mid", "eff"])) + # 事件流: 同 (股, 公告日) 多报告期行罕见(同年同日两期预告)——取最新报告期 + # 为当前信号(P0 语义 = 每公告日一行) + fc_events = (events.sort(["vt_symbol", "eff", "REPORT_DATE"]) + .group_by(["vt_symbol", "eff"]).last() + .select(["vt_symbol", "eff", "forecast_type_score", + "forecast_change_pct", "forecast_np_annualized"])) + return fc_events, fc_pair diff --git a/sanguo_factor/fundamental_pit.py b/sanguo_factor/fundamental_pit.py new file mode 100644 index 0000000..7d7239b --- /dev/null +++ b/sanguo_factor/fundamental_pit.py @@ -0,0 +1,223 @@ +# sanguo_factor/fundamental_pit.py +"""PIT 对齐 + 分块 join 主入口: 报告期/forecast/兑现差/四新域事件流 → +日频特征块(自 fundamental_adapter.py 拆出,纯结构重构). + +R3 PIT 红线: NOTICE_DATE ≤ 决策日才可见(asof backward 前向填充); +按股分批 join_asof 防 NAS 7.9G 内存 OOM. +""" +from __future__ import annotations + +import warnings +from datetime import datetime + +import polars as pl + +from .fundamental_domains import _load_extra_domains +from .fundamental_forecast import _load_forecast_events +from .fundamental_report_features import _compute_report_features, _safe_ratio +from .fundamental_schema import ( + _BEAT_COLS, + _DOMAIN_COLS, + _FORECAST_COLS, + DEFAULT_STATIC_DIR, + FEATURE_COLUMNS, +) +from .fundamental_statements import _load_statements + +# ==================== 对外主入口 ==================== + +# 分块大小: NAS 7.9G 物理内存下,全市场 grid(5555股×2670日×33列≈4G)单次 +# join_asof + 全量副本必 OOM;按股分批使峰值 ≈ 事件表(常驻) + 单批 grid + 输出累积 +BATCH_CODES = 500 + + +def _prepare_days(start: str, end: str, trading_dates) -> list: + """交易日/日历日序列(一次解析,各批共用).""" + if trading_dates is not None: + days = pl.Series("datetime", trading_dates).cast(pl.Datetime("us")) + return days.unique().sort().to_list() + s = datetime.strptime(start, "%Y-%m-%d") + e = datetime.strptime(end, "%Y-%m-%d") + return pl.datetime_range(s, e, interval="1d", eager=True).cast( + pl.Datetime("us")).to_list() + + +def _build_grid(codes: list[str], day_list: list) -> pl.DataFrame: + """单批日频 grid: codes × day_list(批内全量,批间无全量 grid 副本).""" + return pl.DataFrame({ + "vt_symbol": pl.Series([c for c in codes for _ in day_list], dtype=pl.Utf8), + "datetime": pl.Series(day_list * len(codes), dtype=pl.Datetime("us")), + }, schema={"vt_symbol": pl.Utf8, "datetime": pl.Datetime("us")}) + + +def _load_feature_events(codes: list[str], data_dir: str) -> tuple[pl.DataFrame, pl.DataFrame, 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 + and c not in _DOMAIN_COLS] + reports = _load_statements(codes, data_dir) + if reports.height == 0: + stmt_events = pl.DataFrame(schema={ + "vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), **{c: pl.Float64 for c in stmt_cols}}) + beat_events = pl.DataFrame(schema={ + "vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "forecast_beat": pl.Float64}) + else: + feat = _compute_report_features(reports) + # NOTICE_DATE 缺失报告期整期跳过(红线: 宁缺毋假) + feat = feat.filter(pl.col("notice_eff").is_not_null()) + stmt_events = feat.select( + pl.col("vt_symbol"), + pl.col("notice_eff").cast(pl.Datetime("us")).alias("eff"), + *stmt_cols, + ).sort("eff") + beat_events = _build_beat_events(feat, codes, data_dir) + fc_events, _ = _load_forecast_events(codes, data_dir) + fc_events = fc_events.select( + pl.col("vt_symbol"), + pl.col("eff").cast(pl.Datetime("us")), + *_FORECAST_COLS, + ).sort("eff") + return stmt_events, fc_events, beat_events + + +def _build_beat_events(feat: pl.DataFrame, codes: list[str], data_dir: str) -> pl.DataFrame: + """F07 预告兑现差事件流: (实际NP − 预告中值)/abs(预告中值),同 REPORT_DATE 配对. + + 前视红线(P1 任务书): 兑现差含实际 NP,只有实际报告披露后才可知—— + PIT 锚 = max(该报告期 income 有效披露日 notice_eff, 预告公告日), + 不早于两者较晚者(预告公告晚于年报的罕见情形不被提前泄露)。 + 实际 NP 用报告期累计归母净利(预告口径即期间累计);预告中值=0/缺 → NaN。 + """ + schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "forecast_beat": pl.Float64} + _, fc_pair = _load_forecast_events(codes, data_dir) + if fc_pair.height == 0: + return pl.DataFrame(schema=schema) + beat = ( + feat.select("vt_symbol", "REPORT_DATE", "notice_eff", pl.col("np")) + .join(fc_pair, on=["vt_symbol", "REPORT_DATE"], how="inner") + .with_columns( + _safe_ratio(pl.col("np") - pl.col("_fc_mid"), + pl.col("_fc_mid").abs()).alias("forecast_beat"), + pl.max_horizontal("notice_eff", "eff").alias("_anchor"), + ) + .filter(pl.col("forecast_beat").is_not_null() & pl.col("_anchor").is_not_null()) + ) + return beat.select( + pl.col("vt_symbol"), + pl.col("_anchor").cast(pl.Datetime("us")).alias("eff"), + pl.col("forecast_beat"), + ).sort("eff") + + +def iter_fundamental_feature_chunks( + codes: list[str], + start: str, + end: str, + data_dir: str = DEFAULT_STATIC_DIR, + trading_dates: pl.Series | list | None = None, + batch_codes: int = BATCH_CODES, + columns: list[str] | None = None, +): + """按 vt_symbol 分批产出 PIT 日频特征块(生成器,NAS 全量防 OOM 主入口). + + 每块 = 一批 codes × 全部日期 × columns(默认 FEATURE_COLUMNS 全量), + 顺序即 codes 列表顺序;批内 grid 用完即弃,事件右表(报告期+forecast+ + 兑现差)全批共用仅此一份。batch_eval 侧应逐块 join alpha_df 分片后 + concat,避免持有本帧全量副本;columns 子集可只 join 本批表达式引用列, + 全量 68 列 × 1480 万行 ≈ 8G——按引用瘦身是 NAS 7.9G 内存的关键杠杆。 + """ + 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] + grid = _build_grid(chunk_codes, day_list).sort("datetime") + # join_asof 前已显式按键排序;polars 1.42 用 by 分组时无法校验 sortedness, + # 该提示无信息量,就地抑制(sort 即正确性保险) + with warnings.catch_warnings(): + warnings.simplefilter("ignore", UserWarning) + # 三条事件流各自 asof 后丢弃右表键 eff(留置会以 eff_right 后缀 + # 累积,第三次 join 撞名) + out = grid.join_asof( + stmt_events, left_on="datetime", right_on="eff", + by="vt_symbol", strategy="backward").drop("eff") + out = out.sort("datetime").join_asof( + fc_events, left_on="datetime", right_on="eff", + by="vt_symbol", strategy="backward").drop("eff") + 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 # 批间释放(下一批重绑定) + + +def build_fundamental_features( + codes: list[str], + start: str, + end: str, + data_dir: str = DEFAULT_STATIC_DIR, + trading_dates: pl.Series | list | None = None, + batch_codes: int = BATCH_CODES, + columns: list[str] | None = None, +) -> pl.DataFrame: + """构建 PIT 日频财务特征: vt_symbol × datetime × columns. + + Args: + codes: vt_symbol 列表(如 "600000.SSE") + start/end: "YYYY-MM-DD" 窗口(trading_dates=None 时生成日历日 grid) + data_dir: 静态域根目录(NAS=/volume1/stock/sanguo_vnpy_v2/data/static) + trading_dates: 交易日子集(传 bars 的 unique datetime 免造非交易日行) + batch_codes: 按股分批大小(全量防 OOM;测试可调小验分块等值) + columns: 输出特征列子集(默认 FEATURE_COLUMNS 全量;引用瘦身用) + + Returns: + 每行 = 决策日可见的最新报告期特征(NOTICE_DATE ≤ 决策日,asof 前向填充)。 + """ + out_cols = list(columns) if columns is not None else list(FEATURE_COLUMNS) + schema = {"vt_symbol": pl.Utf8, "datetime": pl.Datetime("us"), + **{c: pl.Float64 for c in out_cols}} + if not codes: + return pl.DataFrame(schema=schema) + return pl.concat( + iter_fundamental_feature_chunks( + codes, start, end, data_dir=data_dir, + trading_dates=trading_dates, batch_codes=batch_codes, columns=out_cols), + how="vertical") diff --git a/sanguo_factor/fundamental_report_features.py b/sanguo_factor/fundamental_report_features.py new file mode 100644 index 0000000..0b3f333 --- /dev/null +++ b/sanguo_factor/fundamental_report_features.py @@ -0,0 +1,313 @@ +# sanguo_factor/fundamental_report_features.py +"""报告期级指标预计算: 报告期宽表 → 全部报告期级特征(单季化/TTM/比率族) +(自 fundamental_adapter.py 拆出,纯结构重构). + +R1/R2 口径红线与逐级 with_columns 的装配顺序见 fundamental_adapter.py +模块 docstring;所有 shift/rolling 必须在 .over(_SYM) 组内执行. +""" +from __future__ import annotations + +import polars as pl + +from .fundamental_schema import _CUM_MAP, _DA_PARTS, _FIN_NULL_COLS, _IBD_PARTS, _SYM + +# ==================== 报告期级指标预计算 ==================== +# 所有 shift/rolling 必须在 .over(_SYM) 组内执行(跨股串行 = 致命错误)。 + +def _qidx() -> pl.Expr: + """连续季度索引: year*4 + quarter(月 3/6/9/12 → 1/2/3/4).""" + return (pl.col("REPORT_DATE").dt.year() * 4 + + (pl.col("REPORT_DATE").dt.month() - 1) // 3 + 1) + + +def _single_quarter(col: str) -> pl.Expr: + """R1 单季化: Q1 直接取累计;其余要求上期恰为上一季度(同年)做差,否则 NaN.""" + cur = pl.col(col) + prev = cur.shift(1) + prev_ok = (pl.col("_qidx") - pl.col("_qidx").shift(1)) == 1 + return ( + pl.when(cur.is_null()).then(None) + .when(pl.col("REPORT_DATE").dt.month() == 3).then(cur) + .when(prev_ok & prev.is_not_null()).then(cur - prev) + .otherwise(None) + ).over(_SYM) + + +def _ttm_of(col: str) -> pl.Expr: + """R2 TTM = 连续 4 个报告期单季之和(窗内任一单季 NaN → NaN).""" + q = pl.col(col) + window_ok = (pl.col("_qidx") - pl.col("_qidx").shift(3)) == 3 + s = q + q.shift(1) + q.shift(2) + q.shift(3) + return pl.when(window_ok).then(s).otherwise(None).over(_SYM) + + +def _yoy4(col: str) -> pl.Expr: + """yoy = X_t / X_{t−4季} − 1;基期缺失/≤0 → NaN(负基数 yoy 无意义).""" + cur, base = pl.col(col), pl.col(col).shift(4) + ok = (pl.col("_qidx") - pl.col("_qidx").shift(4)) == 4 + return ( + pl.when(ok & base.is_not_null() & (base > 0) & cur.is_not_null()) + .then(cur / base - 1.0).otherwise(None) + ).over(_SYM) + + +def _delta4(col: str) -> pl.Expr: + """ΔX = X_t − X_{t−4季}(要求恰好隔 4 个季度).""" + ok = (pl.col("_qidx") - pl.col("_qidx").shift(4)) == 4 + return pl.when(ok).then(pl.col(col) - pl.col(col).shift(4)).otherwise(None).over(_SYM) + + +def _mean4(col: str) -> pl.Expr: + """mean(X_t, X_{t−4季})(A03 平均 ROE 分母;恰隔 4 季 + 两期非空守卫).""" + ok = (pl.col("_qidx") - pl.col("_qidx").shift(4)) == 4 + both = pl.col(col).is_not_null() & pl.col(col).shift(4).is_not_null() + m = (pl.col(col) + pl.col(col).shift(4)) / 2.0 + return pl.when(ok & both).then(m).otherwise(None).over(_SYM) + + +def _cagr5(col: str) -> pl.Expr: + """5 年 CAGR = (X_y / X_{y−5})^{1/5} − 1(年报行;基期/现期 ≤0 → NaN, + 负基期 CAGR 无意义;恰隔 20 季守卫).""" + cur, base = pl.col(col), pl.col(col).shift(20) + ok = (pl.col("_qidx") - pl.col("_qidx").shift(20)) == 20 + valid = ok & base.is_not_null() & (base > 0) & cur.is_not_null() & (cur > 0) + return pl.when(valid).then((cur / base).pow(0.2) - 1.0).otherwise(None).over(_SYM) + + +def _std16(col: str) -> pl.Expr: + """16 季滚动 sample std(ddof=1,与 SUE Foster 同款钉死);窗口须恰为连续 + 16 个季度且全非空——不足 16 期/窗内含缺失 → NaN(不填 0).""" + ok = ((pl.col("_qidx") - pl.col("_qidx").shift(15)) == 15).over(_SYM) + sd = pl.col(col).rolling_std(window_size=16, ddof=1).over(_SYM) + return pl.when(ok & sd.is_not_null()).then(sd).otherwise(None) + + +def _safe_ratio(num: pl.Expr, den: pl.Expr) -> pl.Expr: + """分母缺失/为 0 → NaN(比率类通用守卫).""" + return pl.when(den.is_not_null() & (den != 0)).then(num / den).otherwise(None) + + +def _compute_report_features(reports: pl.DataFrame) -> pl.DataFrame: + """报告期宽表 → 全部报告期级特征(逐级 with_columns,over 组内时序).""" + df = reports.sort([_SYM, "REPORT_DATE"]).with_columns(_qidx().alias("_qidx")) + + # 第一级: 单季化 + TTM(累计列) + df = df.with_columns( + [_single_quarter(c).alias(f"q_{c}") for c in _CUM_MAP.values()] + ).with_columns( + [_ttm_of(f"q_{c}").alias(f"ttm_{c}") for c in _CUM_MAP.values()] + ) + + # 第一级半(P1): 补充资料加工底座 + # - DA 累计 = 四件折旧摊销之和(缺列/缺值按 0,与 IBD 组装同款; + # USERIGHT_ASSET_AMORTIZE 2019 前整列缺失不影响早年 DA) + # - 利息费用累计 fill 0(NAS 实测披露稀疏: 600519 78% null/银行模板整列缺; + # 未披露按 0 回加 → EBIT 退化为 TP,口径登记) + # - NWC = 存货+应收(null 传播,缺一即 NaN) + # - 年报行门控列(C13/C14/C18 年度口径;非年报行 null → yoy 天然 NaN) + df = df.with_columns( + pl.sum_horizontal([pl.col(p).fill_null(0.0) for p in _DA_PARTS]).alias("_da_cum"), + pl.col("int_exp").fill_null(0.0).alias("_int0_cum"), + (pl.col("INVENTORY") + pl.col("ACCOUNTS_RECE")).alias("_nwc"), + pl.when(pl.col("REPORT_DATE").dt.month() == 12).then(pl.col("rev")).alias("_rev_ann"), + pl.when(pl.col("REPORT_DATE").dt.month() == 12).then(pl.col("np")).alias("_np_ann"), + pl.when(pl.col("REPORT_DATE").dt.month() == 12).then(pl.col("capex")).alias("_capex_ann"), + ).with_columns( + _single_quarter("_da_cum").alias("q__da"), + _single_quarter("_int0_cum").alias("q__int0"), + ).with_columns( + _ttm_of("q__da").alias("ttm__da"), + _ttm_of("q__int0").alias("ttm__int0"), + ) + + # 第二级: IBD(缺组件按 0,钉死含 LEASE_LIAB 版) + 金融股判定(银行模板无 + # 营业成本) + 毛利 + τ 实际税率(TTM 口径: 消费者均为 TTM 流量; + # τ = INCOME_TAX_TTM/TOTAL_PROFIT_TTM 截断 [0,0.5],两列缺失或 TP≤0 → 0.25) + ibd = pl.sum_horizontal([pl.col(p).fill_null(0.0) for p in _IBD_PARTS]) + is_fin = pl.col("cogs").is_null() | (pl.col("cogs") == 0) + tau = ( + pl.when(pl.col("ttm_tp").is_not_null() & (pl.col("ttm_tp") > 0) + & pl.col("ttm_tax").is_not_null()) + .then((pl.col("ttm_tax") / pl.col("ttm_tp")).clip(0.0, 0.5)) + .otherwise(0.25) + ) + df = df.with_columns( + ibd.alias("_ibd"), + is_fin.alias("_is_fin"), + tau.alias("_tau"), + (pl.col("ttm_rev") - pl.col("ttm_cogs")).alias("_gp_ttm"), + # EBIT_TTM = (TOTAL_PROFIT + FE_INTEREST_EXPENSE)_TTM(§1.3) + (pl.col("ttm_tp") + pl.col("ttm__int0")).alias("_ebit_ttm"), + # (EBIT+DA)_TTM / FCF_TTM / 债务净发行 TTM(E03 三流合成) + (pl.col("ttm_tp") + pl.col("ttm__int0") + pl.col("ttm__da")).alias("_ebitda_ttm"), + (pl.col("ttm_cfo") - pl.col("ttm_capex")).alias("_fcf_ttm"), + (pl.col("ttm_recv_loan") + pl.col("ttm_issue_bond") + - pl.col("ttm_pay_debt")).alias("_debt_issue_ttm"), + ) + + # 第二级半(P1): 平均净资产(A03 分母) + VSIG 三序列底座(单季口径 / TA) + ta = pl.col("TOTAL_ASSETS") + df = df.with_columns( + _mean4("TOTAL_PARENT_EQUITY").alias("_eq_avg"), + _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) + eq = pl.col("TOTAL_PARENT_EQUITY") + df = df.with_columns( + # 盈利能力 A + _safe_ratio(pl.col("ttm_np"), eq).alias("roe_ttm"), + _safe_ratio(pl.col("ttm_dnp"), eq).alias("roe_deduct_ttm"), + _safe_ratio(pl.col("ttm_np"), ta).alias("roa_ttm"), + _safe_ratio(pl.col("_gp_ttm"), ta).alias("gp_over_assets"), + _safe_ratio(pl.col("_gp_ttm"), pl.col("ttm_rev")).alias("gross_margin"), + _safe_ratio(pl.col("ttm_np"), pl.col("ttm_rev")).alias("net_margin"), + _safe_ratio(pl.col("ttm_cfo"), ta).alias("cfo_over_assets"), + # 盈利质量 B + _safe_ratio(pl.col("ttm_np") - pl.col("ttm_cfo"), ta).alias("tacc"), + _safe_ratio(pl.col("ttm_np") - pl.col("ttm_dnp"), + pl.col("ttm_np").abs()).alias("nonrec_ratio"), + _safe_ratio((pl.col("ttm_asset_imp") + pl.col("ttm_credit_imp")).abs(), + ta).alias("impairment_ratio"), + _safe_ratio(pl.col("ttm_invest_inc") + pl.col("ttm_fv_inc"), + pl.col("ttm_tp").abs()).alias("invest_income_dep"), + _safe_ratio(pl.col("ttm_sales_cash"), pl.col("ttm_rev")).alias("sales_cash_ratio"), + _safe_ratio(pl.col("OTHER_RECE"), ta).alias("other_rece_ratio"), + # 成长 C / 资本结构 E + _safe_ratio(pl.col("_ibd"), ta).alias("ibd_ratio"), + _safe_ratio(pl.col("GOODWILL"), ta).alias("goodwill_ratio"), + # 盈利能力 A(P1 7): A03/A09/A10/A11/A12/A15/A16 + _safe_ratio(pl.col("ttm_np"), pl.col("_eq_avg")).alias("roe_avg"), + _safe_ratio(pl.col("ttm_np") + pl.col("ttm__int0") * (1.0 - pl.col("_tau")), + ta).alias("roa_pretax"), + _safe_ratio(pl.col("_ebit_ttm"), ta).alias("ebit_over_assets"), + _safe_ratio(pl.col("_ebitda_ttm"), pl.col("ttm_rev")).alias("ebitda_margin"), + _safe_ratio(pl.col("_ebit_ttm") * (1.0 - pl.col("_tau")), + eq + pl.col("_ibd") - pl.col("MONETARYFUNDS")).alias("roic"), + _safe_ratio(pl.col("ttm_research"), pl.col("ttm_rev")).alias("rd_intensity"), + _safe_ratio(pl.col("ttm_sale_exp"), pl.col("ttm_rev")).alias("sale_expense_ratio"), + # 盈利质量 B(P1): B12 哑变量的连续近似 = (MON/TA)×(IBD/TA) 乘积变体 + # (表达式引擎无截面分位函数,不改引擎——survey B12 的可计算降级) + (_safe_ratio(pl.col("MONETARYFUNDS"), ta) + * _safe_ratio(pl.col("_ibd"), ta)).alias("cash_ibd_product"), + _safe_ratio(pl.col("ttm__da"), pl.col("ttm_rev")).alias("da_intensity"), + # E07 利息保障倍数(利息费用≤0 → NaN: 负利息=净收入,倍数无意义) + 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 组内时序) + df = df.with_columns( + _yoy4("q_rev").alias("rev_q_yoy"), + _yoy4("q_np").alias("np_q_yoy"), + _yoy4("ACCOUNTS_RECE").alias("_ar_yoy"), + _yoy4("ttm_rev").alias("_rev_ttm_yoy"), + _yoy4("SHARE_CAPITAL").alias("nsi"), + _yoy4("TOTAL_ASSETS").alias("asset_growth"), + _delta4("gross_margin").alias("gm_delta"), + _delta4("roe_ttm").alias("roe_delta"), + _delta4("q_np").alias("_diff4_np"), + _delta4("q_rev").alias("_diff4_rev"), + # P1: B09 存货同比 / C16 NWC 同比 / C19 净资产 / E10 商誉 + _yoy4("INVENTORY").alias("_inv_yoy"), + _yoy4("_nwc").alias("nwc_growth"), + _yoy4("TOTAL_PARENT_EQUITY").alias("equity_growth"), + _yoy4("GOODWILL").alias("goodwill_growth"), + # P1: C18 投资增速(年度口径,非年报行 cur=null → NaN) + _yoy4("_capex_ann").alias("invest_growth"), + # P1: C13/C14 五年 CAGR(年报行,恰隔 20 季守卫,基期/现期≤0 → NaN) + _cagr5("_rev_ann").alias("rev_cagr5"), + _cagr5("_np_ann").alias("np_cagr5"), + _delta4("net_margin").alias("nm_delta"), + _delta4("q_eps").alias("_diff4_eps"), + ).with_columns( + (pl.col("np_q_yoy") - pl.col("rev_q_yoy")).alias("growth_scissors"), + (pl.col("_ar_yoy") - pl.col("_rev_ttm_yoy")).alias("receivables_anomaly"), + # B09 与 B08 同构: 期末存量同比 − REV_TTM 同比(登记口径) + (pl.col("_inv_yoy") - pl.col("_rev_ttm_yoy")).alias("inventory_anomaly"), + # B18 毛净剪刀差 = GM_TTM − NM_TTM(第三级产物,同块不可引用故后置) + (pl.col("gross_margin") - pl.col("net_margin")).alias("gm_nm_scissors"), + # C07/C08 加速度 = yoy 的恰隔 4 季二次差分(基期>0 守卫由 yoy 层继承, + # 二次差分同样 NaN 传播;同块不可引用 yoy 列故后置) + _delta4("np_q_yoy").alias("np_accel"), + _delta4("rev_q_yoy").alias("rev_accel"), + ) + + # 第五级: SUE(Foster 标准化)= diff4 / std(过去 8 期 diff4, ddof=1) + for src, out in (("_diff4_np", "sue_np"), ("_diff4_rev", "sue_rev"), + ("_diff4_eps", "sue_eps")): + sd = pl.col(src).rolling_std(window_size=8, ddof=1).over(_SYM) + df = df.with_columns( + pl.when(sd.is_not_null() & (sd > 0) & pl.col(src).is_not_null()) + .then(pl.col(src) / sd).otherwise(None).alias(out) + ) + # SUE 严窗变体(随批互评): σ 只用 t−1 及更早差分(shift(1) 后滚 8 期,不含当期) + sd_strict = pl.col("_diff4_np").shift(1).rolling_std(window_size=8, ddof=1).over(_SYM) + df = df.with_columns( + pl.when(sd_strict.is_not_null() & (sd_strict > 0) & pl.col("_diff4_np").is_not_null()) + .then(pl.col("_diff4_np") / sd_strict).otherwise(None).alias("sue_np_strict") + ) + + # 第五级半(P1): VSIG 16 季滚动(sample std ddof=1 钉死,恰连续 16 季全非空) + # + F06 披露及时性 = −(有效披露日 − 报告期) 天数(早披露=高分;notice_eff + # 取三表最晚可见,与 PIT 锚一致) + df = df.with_columns( + _std16("_np_ta").alias("vsig"), + _std16("_accq_ta").alias("vsig_acc"), + _std16("_cfo_ta").alias("vsig_cfo"), + (-(pl.col("notice_eff") - pl.col("REPORT_DATE")).dt.total_days()) + .cast(pl.Float64).alias("disclosure_speed"), + ) + + # 第五级半续(P1): B19 持续盈利季数 = 连续单季 NP>0 计数(截断 8; + # 当期缺失→NaN;非正→0 断流;中间缺失行视为断流点) + df = df.with_columns( + (pl.col("q_np") > 0).alias("_pos"), + pl.int_range(pl.len()).cast(pl.Int64).alias("_ridx"), + ).with_columns( + pl.when(pl.col("_pos").is_null() | ~pl.col("_pos")) + .then(pl.col("_ridx")).otherwise(None) + .fill_null(strategy="forward").over(_SYM).alias("_lastbrk"), + ).with_columns( + pl.when(pl.col("_pos").is_null()).then(None) + .when(pl.col("_pos")) + .then((pl.col("_ridx") - pl.col("_lastbrk").fill_null(-1)).clip(1, 8).cast(pl.Float64)) + .otherwise(0.0).alias("profit_streak"), + ) + + # 金融股: 盈利质量/成长/费用类族特征置 NaN(§7 红线 5;估值族保留) + df = df.with_columns([ + pl.when(pl.col("_is_fin")).then(None).otherwise(pl.col(c)).alias(c) + for c in _FIN_NULL_COLS + ]) + + # 输出别名(估值/资本行为因子的日频原料) + df = df.with_columns( + pl.col("ttm_np").alias("np_ttm"), + pl.col("ttm_dnp").alias("dnp_ttm"), + pl.col("ttm_cfo").alias("cfo_ttm"), + pl.col("ttm_acc_inv_cash").alias("acc_invest_cash_ttm"), + pl.col("TOTAL_PARENT_EQUITY").alias("equity"), + pl.col("SHARE_CAPITAL").alias("share_capital"), + # P1: EV 群/FCF/债务净发行原料(EV = close×share_capital + ev_ex_mv) + pl.col("_gp_ttm").alias("gp_ttm"), + pl.col("_ebit_ttm").alias("ebit_ttm"), + pl.col("_ebitda_ttm").alias("ebitda_ttm"), + pl.col("_fcf_ttm").alias("fcf_ttm"), + (pl.col("_ibd") - pl.col("MONETARYFUNDS")).alias("ev_ex_mv"), + pl.col("_debt_issue_ttm").alias("debt_issue_ttm"), + ) + return df diff --git a/sanguo_factor/fundamental_schema.py b/sanguo_factor/fundamental_schema.py new file mode 100644 index 0000000..556af03 --- /dev/null +++ b/sanguo_factor/fundamental_schema.py @@ -0,0 +1,128 @@ +# sanguo_factor/fundamental_schema.py +"""财务因子适配层常量/列契约(自 fundamental_adapter.py 拆出,纯结构重构). + +口径红线与数据形态的完整说明见 fundamental_adapter.py 模块 docstring; +本模块只承载各数据域共享的列映射/清单常量,被读取层/特征层/PIT 层共同引用. +""" +from __future__ import annotations + +DEFAULT_STATIC_DIR = "/volume1/stock/sanguo_vnpy_v2/data/static" +_SYM = "vt_symbol" # 时序分组键 + +# vt_symbol 后缀 → NAS 文件名后缀(600000.SSE → 600000.SH) +_VT_TO_FILE_SUFFIX = {"SSE": "SH", "SZSE": "SZ"} + +# forecast 预告类型 → 有序分(§3.6 F04) +FORECAST_TYPE_SCORE: dict[str, int] = { + "预增": 3, "略增": 2, "扭亏": 2, "续盈": 1, "减亏": 1, "不确定": 0, + "略减": -1, "增亏": -2, "续亏": -2, "预减": -3, "首亏": -3, +} + +# 累计口径列映射: 源列 → 短名(需单季化+TTM) +_CUM_MAP = { + "TOTAL_OPERATE_INCOME": "rev", + "OPERATE_COST": "cogs", + "PARENT_NETPROFIT": "np", + "DEDUCT_PARENT_NETPROFIT": "dnp", + "TOTAL_PROFIT": "tp", + "INVEST_INCOME": "invest_inc", + "FAIRVALUE_CHANGE_INCOME": "fv_inc", + "ASSET_IMPAIRMENT_LOSS": "asset_imp", + "CREDIT_IMPAIRMENT_LOSS": "credit_imp", + "NETCASH_OPERATE": "cfo", + "SALES_SERVICES": "sales_cash", + "ACCEPT_INVEST_CASH": "acc_inv_cash", + # ---- P1 批新增(income;NAS 实测 2026-09-08 列名核实)---- + "RESEARCH_EXPENSE": "research", # 研发费用(2018Q3 起单列,早年 null 传播) + "SALE_EXPENSE": "sale_exp", # 销售费用 + "FE_INTEREST_EXPENSE": "int_exp", # 财务费用-利息费用(披露稀疏,EBIT 组装层按 0) + "INCOME_TAX": "tax", # 所得税费用(τ 实际税率分子) + "BASIC_EPS": "eps", # 基本每股收益(累计口径,需单季化) + # ---- P1 批新增(cashflow 补充资料段 + 融资流)---- + "FA_IR_DEPR": "fa_depr", # 固定资产折旧 + "IA_AMORTIZE": "ia_amort", # 无形资产摊销 + "LPE_AMORTIZE": "lpe_amort", # 长期待摊费用摊销 + "USERIGHT_ASSET_AMORTIZE": "ua_amort", # 使用权资产折旧摊销(2019 起才有) + "CONSTRUCT_LONG_ASSET": "capex", # 购建长期资产现金(C17 投资率/D15 FCF) + "RECEIVE_LOAN_CASH": "recv_loan", # 取得借款现金(E03) + "ISSUE_BOND": "issue_bond", # 发行债券现金(E03) + "PAY_DEBT_CASH": "pay_debt", # 偿还债务现金(E03) +} +# 存量列(balance,时点值直接用;缺列 → null) +_BALANCE_COLS = ["TOTAL_ASSETS", "TOTAL_PARENT_EQUITY", "ACCOUNTS_RECE", + "OTHER_RECE", "GOODWILL", "SHARE_CAPITAL", "SHORT_LOAN", + "SHORT_FIN_PAYABLE", "NONCURRENT_LIAB_1YEAR", "LONG_LOAN", + "BOND_PAYABLE", "LEASE_LIAB", + # P1 新增存量列 + "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", + "LONG_LOAN", "BOND_PAYABLE", "LEASE_LIAB"] +# DA 组装: cashflow 间接法补充资料四件,列/值缺失按 0(「列存在才加」; +# USERIGHT_ASSET_AMORTIZE 2019 前整列缺失不影响早年 DA) +_DA_PARTS = ["fa_depr", "ia_amort", "lpe_amort", "ua_amort"] +_DATE_COLS = ["REPORT_DATE", "NOTICE_DATE", "UPDATE_DATE"] + +# 输出特征列(P0 32 + P1 35 因子及 2 变体的全部原料;契约由 test_fundamental_library +# 与 test_fundamental_p1_library 锁定) +FEATURE_COLUMNS: list[str] = [ + # 报告期级比率(盈利能力 A / 盈利质量 B / 成长 C / 资本结构 E / 预期事件 F) + "roe_ttm", "roe_deduct_ttm", "roa_ttm", "gp_over_assets", "gross_margin", + "net_margin", "cfo_over_assets", + "tacc", "nonrec_ratio", "impairment_ratio", "invest_income_dep", + "receivables_anomaly", "sales_cash_ratio", "other_rece_ratio", + "rev_q_yoy", "np_q_yoy", "growth_scissors", "gm_delta", "roe_delta", + "asset_growth", "nsi", "ibd_ratio", "goodwill_ratio", + "sue_np", "sue_rev", + "forecast_type_score", "forecast_change_pct", + # 估值/资本行为因子的日频原料(表达式层 ÷ close×share_capital) + "np_ttm", "dnp_ttm", "cfo_ttm", "equity", "share_capital", + "acc_invest_cash_ttm", + # ---- P1 批新增(A 盈利 7)---- + "roe_avg", "roa_pretax", "ebit_over_assets", "ebitda_margin", "roic", + "rd_intensity", "sale_expense_ratio", + # ---- P1 批新增(B 质量 8;B12 为哑变量的连续近似乘积)---- + "inventory_anomaly", "cash_ibd_product", "vsig", "vsig_acc", "vsig_cfo", + "da_intensity", "gm_nm_scissors", "profit_streak", + # ---- P1 批新增(C 成长 8)---- + "np_accel", "rev_accel", "nm_delta", "rev_cagr5", "np_cagr5", + "nwc_growth", "invest_growth", "equity_growth", + # ---- P1 批新增(D 估值 6 的原料;EV = close×share_capital + ev_ex_mv; + # forecast_np_annualized 走 forecast 事件流公告日 asof)---- + "ebit_ttm", "ebitda_ttm", "fcf_ttm", "gp_ttm", "ev_ex_mv", + "forecast_np_annualized", + # ---- P1 批新增(E 资本 3)---- + "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 同剔) +_FIN_NULL_COLS = ["tacc", "nonrec_ratio", "impairment_ratio", + "invest_income_dep", "receivables_anomaly", + "sales_cash_ratio", "other_rece_ratio", + "rev_q_yoy", "np_q_yoy", "growth_scissors", "gm_delta", + "roe_delta", "asset_growth", + "rd_intensity", "sale_expense_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", + # 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") diff --git a/sanguo_factor/fundamental_statements.py b/sanguo_factor/fundamental_statements.py new file mode 100644 index 0000000..4c710a9 --- /dev/null +++ b/sanguo_factor/fundamental_statements.py @@ -0,0 +1,130 @@ +# sanguo_factor/fundamental_statements.py +"""三表读取层: NAS static 域 income/balance/cashflow parquet → 报告期宽表 +(自 fundamental_adapter.py 拆出,纯结构重构). + +按股文件路由 + 容错跳过(零行/缺文件/坏文件)+ (vt_symbol, REPORT_DATE) +去重取重述终值,有效披露日 notice_eff = 三表 NOTICE_DATE 最大值(保守). +""" +from __future__ import annotations + +import os + +import polars as pl + +from .fundamental_schema import _BALANCE_COLS, _CUM_MAP, _DATE_COLS, _VT_TO_FILE_SUFFIX + +# ==================== 读取层 ==================== + +def _vt_to_file_code(vt_symbol: str) -> str | None: + """``600000.SSE`` → ``600000.SH``;无法映射的(ETF/北交)返回 None.""" + code, _, suffix = vt_symbol.partition(".") + file_suffix = _VT_TO_FILE_SUFFIX.get(suffix.upper()) + return f"{code}.{file_suffix}" if file_suffix else None + + +def _read_static(path: str, want_cols: list[str]) -> pl.DataFrame | None: + """读单股单表 parquet,缺列补 null;零行/缺文件/坏文件 → None(容错跳过).""" + if not os.path.exists(path): + return None + try: + schema = pl.read_parquet_schema(path) + if "REPORT_DATE" not in schema: + return None + cols = [c for c in _DATE_COLS + want_cols if c in schema] + df = pl.read_parquet(path, columns=cols) + except Exception: + return None + if df.height == 0: + return None + missing = [c for c in _DATE_COLS + want_cols if c not in df.columns] + if missing: + df = df.with_columns([pl.lit(None, dtype=pl.Utf8).alias(c) for c in missing]) + for c in want_cols: + df = df.with_columns(pl.col(c).cast(pl.Float64, strict=False)) + return df + + +def _norm_dates(df: pl.DataFrame, cols: list[str]) -> pl.DataFrame: + """日期列归一为 pl.Date:字符串截前 10 位解析,Date/Datetime 直接 cast.""" + exprs = [] + for c in cols: + dtype = df.schema[c] + if dtype == pl.Utf8: + exprs.append(pl.col(c).str.slice(0, 10).str.to_date("%Y-%m-%d", strict=False).alias(c)) + elif dtype != pl.Date: + exprs.append(pl.col(c).cast(pl.Date, strict=False).alias(c)) + return df.with_columns(exprs) if exprs else df + + +def _dedupe_reports(df: pl.DataFrame, value_cols: list[str]) -> pl.DataFrame: + """按 (vt_symbol, REPORT_DATE) 去重:值取 (UPDATE_DATE, NOTICE_DATE) 排序末行 + (重述取终值),有效披露日取组内 NOTICE_DATE 最大值(保守,不提前).""" + df = df.sort(["vt_symbol", "REPORT_DATE", "UPDATE_DATE", "NOTICE_DATE"], + nulls_last=False) + aggs = [pl.col(c).last().alias(c) for c in value_cols] + aggs.append(pl.col("NOTICE_DATE").max().alias("notice_eff")) + return df.group_by(["vt_symbol", "REPORT_DATE"]).agg(aggs) + + +# 各表显式读取清单(income/cashflow 列集不相交,join 不加后缀) +_TABLE_RAW = { + "income": ["TOTAL_OPERATE_INCOME", "OPERATE_COST", "PARENT_NETPROFIT", + "DEDUCT_PARENT_NETPROFIT", "TOTAL_PROFIT", "INVEST_INCOME", + "FAIRVALUE_CHANGE_INCOME", "ASSET_IMPAIRMENT_LOSS", + "CREDIT_IMPAIRMENT_LOSS", + "RESEARCH_EXPENSE", "SALE_EXPENSE", "FE_INTEREST_EXPENSE", + "INCOME_TAX", "BASIC_EPS"], + "balance": _BALANCE_COLS, + "cashflow": ["NETCASH_OPERATE", "SALES_SERVICES", "ACCEPT_INVEST_CASH", + "FA_IR_DEPR", "IA_AMORTIZE", "LPE_AMORTIZE", + "USERIGHT_ASSET_AMORTIZE", "CONSTRUCT_LONG_ASSET", + "RECEIVE_LOAN_CASH", "ISSUE_BOND", "PAY_DEBT_CASH"], +} + + +def _load_statements(codes: list[str], data_dir: str) -> pl.DataFrame: + """全 codes 三表 → 报告期宽表(每 vt_symbol × REPORT_DATE 一行). + + income/cashflow 值列重命名为短名(_CUM_MAP),balance 保持源列名; + anchor = 三表报告期并集,left join 保证单表缺期不拖垮其它表特征。 + """ + table_cols = _TABLE_RAW + per_table: dict[str, pl.DataFrame] = {} + for table, raw_cols in table_cols.items(): + frames = [] + for vt in codes: + file_code = _vt_to_file_code(vt) + if file_code is None: + continue + df = _read_static( + os.path.join(data_dir, table, f"{file_code}_{table}.parquet"), raw_cols) + if df is None: + continue + df = _norm_dates(df, _DATE_COLS) + # 先统一列序/列集再入列(各股文件 schema 子集不同,concat 前必须对齐) + frames.append( + df.with_columns(pl.lit(vt).alias("vt_symbol")) + .select(["vt_symbol", *_DATE_COLS, *raw_cols])) + if not frames: + continue + merged = pl.concat(frames) + if table == "balance": + short = {c: c for c in raw_cols} + else: + short = {k: v for k, v in _CUM_MAP.items() if k in raw_cols} + dedup = _dedupe_reports(merged, list(raw_cols)).rename(short) + per_table[table] = dedup + + empty = pl.DataFrame(schema={"vt_symbol": pl.Utf8, "REPORT_DATE": pl.Date}) + if not per_table: + return empty + + out = pl.concat([t.select(["vt_symbol", "REPORT_DATE"]) + for t in per_table.values()]).unique() + notice_cols = [] + for table, dedup in per_table.items(): + renamed = dedup.rename({"notice_eff": f"_notice_{table}"}) + notice_cols.append(f"_notice_{table}") + out = out.join(renamed, on=["vt_symbol", "REPORT_DATE"], how="left") + # 有效披露日 = 三表 NOTICE_DATE 行最大(最晚可见,保守不提前) + return out.with_columns(pl.max_horizontal(notice_cols).alias("notice_eff"))