Files
sanguo_vnpy_v2/sanguo_factor/fundamental_adapter.py
T
claude_dev e330e130d4 chore(factor): P1二轮回评修缮——v1.1口径残留对齐+导入示例纠偏+replace_strict [nas]
- composite_library docstring/测试注释: 7源/19源/fund7/all19 残留 → 6源/18源
  (7dfcaca 剔 topholder 后文档未同步,防后续按旧口径引用)
- import_runs_to_main_db 示例 --dst data_backup→data(09-08 导错库坑,
  示例即陷阱原文,加⚠️一行防再踩)
- forecast 预告类型映射 replace→replace_strict(消 398 条 DeprecationWarning,
  参数语义等价,231 测全绿)
- P1 晨报§五: 补 09-09 只读探针复核(首期 20210930/5431文件每期)+§19.10
  修订文本入档(master 侧一行改,本分支无该节)

Co-Authored-By: Claude Code <noreply@anthropic.com>
2026-09-09 12:12:00 +08:00

1202 lines
59 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# sanguo_factor/fundamental_adapter.py
"""财务因子适配层: NAS 静态域三表+forecast parquet → PIT 日频特征列.
口径红线(docs/fundamental_factor_survey_20260907.md §7):
- R1 单季差分: Q1 直接取累计,其余 = 本期累计 − 年内上期累计;缺上期置 NaN 不填 0
- R2 TTM: 连续 4 个单季之和,不足 4 期 NaN
- R3 PIT: NOTICE_DATE ≤ 决策日才可见(有效披露日 = 三表 NOTICE_DATE 最大值,保守);
NOTICE_DATE 缺失的报告期整期跳过(宁缺毋假)
- 金融股(OPERATE_COST 缺失/为 0 的银行模板): 盈利质量/成长族特征置 NaN(估值族保留)
数据形态(NAS 实测 2026-09-07):
- 三表按股: static/{income,balance,cashflow}/{code}.{SH|SZ}_{table}.parquet
(北交 920xxx 残留零行文件、沪深退市 355 只文件不存在 → 读取容错跳过)
- forecast 按报告期全市场: static/forecast/{YYYYMMDD}_forecast.parquet(中文列,
一股多行=按「预测指标」,归母净利润行优先)
- 日期列为 "YYYY-MM-DD 00:00:00" 字符串;东财金额单位=元
- valuation(中文列)不读: 估值类因子市值 = close × SHARE_CAPITAL 自算(表达式层)
P1-B 四新域(NAS 实测 2026-09-08):
- dividend 按股: static/dividend/{code}.{SH|SZ}_dividend.parquet;code 列为
baostock 小写格式(sh.600519)与文件名不同 → 以文件名路由;日期/数值列为
字符串('' 为空);一行=一次分红事件,同年可多行
- gdhs 按期: static/gdhs/{YYYYMMDD}_gdhs.parquet(53 期 2013Q1→2026Q1);
用 代码/股东户数-本次/股东户数统计截止日-本次/公告日期 四列(实测列名
无第二连字符;噪声行情快照列不读)
- top_holders: 只读聚合产物 factor_cache/top_holders_agg.parquet(static
兄弟目录;由 scripts/factor_research/preaggregate_top_holders.py 一次性
预聚合 108,610 个按股×期小文件)
- valuation_baostock 按年: data/valuation_baostock/{year}.parquet(static
兄弟目录);date 为字符串,exchange 实测为 SH/SZ(映射回 SSE/SZSE);
2026 文件仅 08-13 起(bs 日喂起点)→ 01-01~08-12
缺口用 static/valuation(ak em 中文列 PE(TTM)/市净率)倒数补
- 四域缺任一 → 相关特征列全 NaN + 一次性 warning(优雅降级,本地可跑通)
输出: vt_symbol × datetime(日频) × FEATURE_COLUMNS,join 到 alpha_df 后
供 cs_rank(列) 表达式直接消费。
"""
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_{t4季}(要求恰好隔 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_{t4季})(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_{y5})^{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+DIODPO(缺任一原料 NaN 传播)
(_safe_ratio(pl.col("_ar_avg") * 365.0, pl.col("ttm_rev"))
+ _safe_ratio(pl.col("_inv_avg") * 365.0, pl.col("ttm_cogs"))
- _safe_ratio(pl.col("_ap_avg") * 365.0, pl.col("ttm_cogs"))).alias("ccc"),
)
# 第四级: 跨期差分/同比/剪刀差(over 组内时序)
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(d365) = 窗 (d365, d] 内事件合计)。
code 列为 baostock 小写格式(sh.600519)与文件名不同 → 以文件名路由
(与三表读取同模式,列内容不参与 join)。
"""
schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "_send_cum": pl.Float64}
div_dir = os.path.join(data_dir, "dividend")
if not os.path.isdir(div_dir):
_warn_domain_missing("dividend", div_dir)
return pl.DataFrame(schema=schema)
frames = []
for vt in codes:
file_code = _vt_to_file_code(vt)
if file_code is None:
continue
path = os.path.join(div_dir, f"{file_code}_dividend.parquet")
if not os.path.exists(path):
continue
try:
df = pl.read_parquet(path)
except Exception:
continue
chain = [c for c in _DIVIDEND_CHAIN if c in df.columns]
if not chain or "dividStocksPs" not in df.columns:
continue
# 日期列字符串('' 为空;兼容 'YYYY-MM-DD 00:00:00')→ coalesce 回退链
eff = pl.coalesce([
pl.col(c).cast(pl.Utf8).str.slice(0, 10).str.to_date("%Y-%m-%d", strict=False)
for c in chain])
df = df.select(
pl.lit(vt).alias("vt_symbol"),
eff.alias("eff"),
pl.col("dividStocksPs").cast(pl.Float64, strict=False).alias("_send"),
).filter(pl.col("eff").is_not_null())
if df.height:
frames.append(df)
if not frames:
return pl.DataFrame(schema=schema)
ev = pl.concat(frames).sort(["vt_symbol", "eff"])
return ev.with_columns(
pl.col("_send").cum_sum().over(_SYM).alias("_send_cum")
).select(["vt_symbol", "eff", "_send_cum"]).with_columns(
pl.col("eff").cast(pl.Datetime("us")))
def _load_gdhs_events(codes: list[str], data_dir: str) -> pl.DataFrame:
"""gdhs 按期文件 → (vt_symbol, eff, gdhs_chg) 事件流.
事件表: (代码, 统计截止日)→股东户数;同键多行取最新公告(修订口径);
gdhs_chg = 相邻事件(按截止日排序)Δln(户数),首事件无上期 → NaN;
PIT = 公告日期。红线: 不用「股东户数-增减比例」列(每股截止日不规则,
非统一季环比)。
"""
schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "gdhs_chg": pl.Float64}
gd_dir = os.path.join(data_dir, "gdhs")
if not os.path.isdir(gd_dir):
_warn_domain_missing("gdhs", gd_dir)
return pl.DataFrame(schema=schema)
code_set = set(codes)
frames = []
for fname in sorted(os.listdir(gd_dir)):
if not fname.endswith(".parquet"):
continue
try:
f = pl.read_parquet(os.path.join(gd_dir, fname))
except Exception:
continue
cutoff = next((c for c in _GDHS_CUTOFF_CANDIDATES if c in f.columns), None)
if cutoff is None or not all(
c in f.columns for c in ("代码", "股东户数-本次", "公告日期")):
continue
code = pl.col("代码").cast(pl.Utf8).str.strip_chars().str.zfill(6)
f = f.select(
_code6_to_vt(code).alias("vt_symbol"),
pl.col("股东户数-本次").cast(pl.Float64, strict=False).alias("_holders"),
pl.col(cutoff).cast(pl.Date, strict=False).alias("_cutoff"),
pl.col("公告日期").cast(pl.Date, strict=False).alias("_ann"),
).filter(pl.col("vt_symbol").is_in(code_set) & pl.col("_ann").is_not_null()
& pl.col("_cutoff").is_not_null() & pl.col("_holders").is_not_null())
if f.height:
frames.append(f)
if not frames:
return pl.DataFrame(schema=schema)
ev = (pl.concat(frames)
.sort(["vt_symbol", "_cutoff", "_ann"])
.group_by(["vt_symbol", "_cutoff"]).agg( # 同 (股, 截止日) 取最新公告
pl.col("_holders").last(), pl.col("_ann").last())
.sort(["vt_symbol", "_cutoff"]))
return ev.with_columns(
pl.col("_holders").shift(1).over(_SYM).alias("_prev")
).with_columns(
pl.when((pl.col("_prev") > 0) & (pl.col("_holders") > 0))
.then((pl.col("_holders") / pl.col("_prev")).log())
.otherwise(None).alias("gdhs_chg")
).select(
pl.col("vt_symbol"),
pl.col("_ann").cast(pl.Datetime("us")).alias("eff"),
pl.col("gdhs_chg"),
).sort("eff")
def _load_topholder_events(codes: list[str], data_dir: str) -> pl.DataFrame:
"""top_holders 聚合产物 → (vt_symbol, eff, topholder_chg) 事件流.
聚合产物 = scripts/factor_research/preaggregate_top_holders.py 输出
(file_code, period, hold_pct, n_holders),置于 static 兄弟目录
factor_cache/ 下(绝不写 static 树)。期→PIT 无法定披露日 → 法定披露
截止近似(保守侧): Q1→04-30 / H1→08-31 / Q3→10-31 / 年报→次年04-30。
topholder_chg = 最近期十大合计占比 − 上期(首期无上期 → NaN)。
"""
schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"), "topholder_chg": pl.Float64}
agg_path = os.path.join(os.path.dirname(os.path.abspath(data_dir)),
"factor_cache", "top_holders_agg.parquet")
if not os.path.exists(agg_path):
_warn_domain_missing("top_holders_agg", agg_path)
return pl.DataFrame(schema=schema)
try:
agg = pl.read_parquet(agg_path)
except Exception:
_warn_domain_missing("top_holders_agg", agg_path)
return pl.DataFrame(schema=schema)
if not all(c in agg.columns for c in ("file_code", "period", "hold_pct")):
return pl.DataFrame(schema=schema)
parts = pl.col("file_code").str.split(".")
ev = agg.with_columns(
parts.list.get(0).alias("_c"), parts.list.get(1).alias("_x"),
pl.col("period").cast(pl.Date, strict=False),
pl.col("hold_pct").cast(pl.Float64, strict=False),
).with_columns(
pl.when(pl.col("_x") == "SH").then(pl.col("_c") + pl.lit(".SSE"))
.otherwise(pl.col("_c") + pl.lit(".SZSE")).alias("vt_symbol")
).filter(
pl.col("vt_symbol").is_in(set(codes)) & pl.col("period").is_not_null()
& pl.col("hold_pct").is_not_null()
).sort(["vt_symbol", "period"])
m = pl.col("period").dt.month()
y = pl.col("period").dt.year()
deadline = (
pl.when(m == 3).then(pl.date(y, 4, 30))
.when(m == 6).then(pl.date(y, 8, 31))
.when(m == 9).then(pl.date(y, 10, 31))
.otherwise(pl.date(y + 1, 4, 30))
)
return ev.with_columns(
pl.col("hold_pct").shift(1).over(_SYM).alias("_prev")
).with_columns(
(pl.col("hold_pct") - pl.col("_prev")).alias("topholder_chg")
).select(
pl.col("vt_symbol"),
deadline.cast(pl.Datetime("us")).alias("eff"),
pl.col("topholder_chg"),
).sort("eff")
def _inv_guard(col: str) -> pl.Expr:
"""1/x 守卫: x 缺失/为 0 → NaN(亏损 baostock pe=NaN 天然 NaN)."""
return (pl.when(pl.col(col).is_not_null() & (pl.col(col) != 0))
.then(1.0 / pl.col(col)).otherwise(None))
def _load_vb_daily(codes: list[str], data_dir: str,
start: str, end: str) -> pl.DataFrame:
"""valuation_baostock 按年文件(+static/valuation 补 2026 缺口) →
(vt_symbol, eff, ep_vb, bp_vb) 日频流.
vb 在 static 的兄弟目录(data/valuation_baostock/{year}.parquet),
date 为字符串须转 Date;ep=1/peTTM、bp=1/pbMRQ。2026-01-01~08-12
缺口(bs 日喂 08-13 起)用 static/valuation(ak em 中文列:
PE(TTM)/市净率)倒数补——两源亏损口径不同(baostock NaN/ak 可为负),
倒数后均按原始符号保留,分域验证时注意。
"""
schema = {"vt_symbol": pl.Utf8, "eff": pl.Datetime("us"),
"ep_vb": pl.Float64, "bp_vb": pl.Float64}
root = os.path.dirname(os.path.abspath(data_dir))
vb_dir = os.path.join(root, "valuation_baostock")
if not os.path.isdir(vb_dir):
_warn_domain_missing("valuation_baostock", vb_dir)
return pl.DataFrame(schema=schema)
code_set = set(codes)
d_lo = datetime.strptime(start, "%Y-%m-%d").date()
d_hi = datetime.strptime(end, "%Y-%m-%d").date()
frames = []
for year in range(d_lo.year, d_hi.year + 1):
path = os.path.join(vb_dir, f"{year}.parquet")
if not os.path.exists(path):
continue
try:
f = pl.read_parquet(path)
except Exception:
continue
if not all(c in f.columns for c in
("symbol", "exchange", "date", "peTTM", "pbMRQ")):
continue
raw_date = pl.col("date")
d = (raw_date.cast(pl.Utf8).str.slice(0, 10).str.to_date("%Y-%m-%d", strict=False)
if f.schema["date"] == pl.Utf8 else raw_date.cast(pl.Date, strict=False))
# exchange 实测为 SH/SZ(NAS 2026-09-08,任务书 SSE/SZSE 变体兼容)
xsuf = (pl.when(pl.col("exchange") == pl.lit("SH")).then(pl.lit("SSE"))
.when(pl.col("exchange") == pl.lit("SZ")).then(pl.lit("SZSE"))
.otherwise(pl.col("exchange")))
f = f.select(
(pl.col("symbol").cast(pl.Utf8).str.strip_chars() + pl.lit(".")
+ xsuf.cast(pl.Utf8).str.strip_chars()).alias("vt_symbol"),
d.alias("eff"),
pl.col("peTTM").cast(pl.Float64, strict=False).alias("_pe"),
pl.col("pbMRQ").cast(pl.Float64, strict=False).alias("_pb"),
).filter(pl.col("vt_symbol").is_in(code_set) & pl.col("eff").is_not_null()
& (pl.col("eff") >= d_lo) & (pl.col("eff") <= d_hi))
if f.height:
frames.append(f)
# 2026 缺口补口: static/valuation 按股中文列
gap_lo, gap_hi = max(_VB_GAP[0], d_lo), min(_VB_GAP[1], d_hi)
if gap_lo <= gap_hi:
val_dir = os.path.join(data_dir, "valuation")
if not os.path.isdir(val_dir):
_warn_domain_missing("static/valuation(2026 补口)", val_dir)
else:
for vt in codes:
file_code = _vt_to_file_code(vt)
if file_code is None:
continue
path = os.path.join(val_dir, f"{file_code}_valuation.parquet")
if not os.path.exists(path):
continue
try:
f = pl.read_parquet(path)
except Exception:
continue
have = [c for c in ("数据日期", "PE(TTM)", "市净率") if c in f.columns]
if "数据日期" not in have:
continue
f = _norm_dates(f.select(have), ["数据日期"])
pe = (pl.col("PE(TTM)").cast(pl.Float64, strict=False)
if "PE(TTM)" in have else pl.lit(None, pl.Float64))
pb = (pl.col("市净率").cast(pl.Float64, strict=False)
if "市净率" in have else pl.lit(None, pl.Float64))
f = f.select(
pl.lit(vt).alias("vt_symbol"),
pl.col("数据日期").alias("eff"),
pe.alias("_pe"), pb.alias("_pb"),
).filter(pl.col("eff").is_not_null()
& (pl.col("eff") >= gap_lo) & (pl.col("eff") <= gap_hi))
if f.height:
frames.append(f)
if not frames:
return pl.DataFrame(schema=schema)
out = (pl.concat(frames).sort(["vt_symbol", "eff"])
.unique(subset=["vt_symbol", "eff"], keep="last")
.with_columns(
_inv_guard("_pe").alias("ep_vb"),
_inv_guard("_pb").alias("bp_vb"),
pl.col("eff").cast(pl.Datetime("us"))))
return out.select(["vt_symbol", "eff", "ep_vb", "bp_vb"]).sort("eff")
_EMPTY_EVENTS = pl.DataFrame(schema={"vt_symbol": pl.Utf8, "eff": pl.Datetime("us")})
def _load_extra_domains(codes: list[str], data_dir: str,
start: str, end: str, out_cols: list[str]) -> dict:
"""P1-B 四新域按引用列加载(引用瘦身的域级延伸: 未引用的域不读盘)."""
cols = set(out_cols)
return {
"dividend": (_load_dividend_cum(codes, data_dir)
if "send_total_12m" in cols else _EMPTY_EVENTS),
"gdhs": (_load_gdhs_events(codes, data_dir)
if "gdhs_chg" in cols else _EMPTY_EVENTS),
"topholder": (_load_topholder_events(codes, data_dir)
if "topholder_chg" in cols else _EMPTY_EVENTS),
"vb": (_load_vb_daily(codes, data_dir, start, end)
if "ep_vb" in cols or "bp_vb" in cols else _EMPTY_EVENTS),
}
# ==================== 对外主入口 ====================
# 分块大小: NAS 7.9G 物理内存下,全市场 grid(5555股×2670日×33列≈4G)单次
# 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(d365) = 窗 (d365, d] 事件合计(开区间下界);
# 域内无事件(≤d) → NaN 不填 0(与缺文件/域不覆盖区分)
div365 = div.rename({"_send_cum": "_send_cum_365"})
out = out.sort("datetime").join_asof(
div, left_on="datetime", right_on="eff",
by="vt_symbol", strategy="backward").drop("eff")
out = (
out.with_columns(
(pl.col("datetime") - pl.duration(days=365)).alias("_dt365"))
.sort("_dt365")
.join_asof(div365, left_on="_dt365", right_on="eff",
by="vt_symbol", strategy="backward")
.drop("eff", "_dt365")
.with_columns(
pl.when(pl.col("_send_cum").is_null()).then(None)
.otherwise(pl.col("_send_cum")
- pl.col("_send_cum_365").fill_null(0.0))
.alias("send_total_12m")))
for key in ("gdhs", "topholder"):
if extra[key].height:
out = out.sort("datetime").join_asof(
extra[key], left_on="datetime", right_on="eff",
by="vt_symbol", strategy="backward").drop("eff")
if extra["vb"].height:
out = out.sort("datetime").join_asof(
extra["vb"], left_on="datetime", right_on="eff",
by="vt_symbol", strategy="backward").drop("eff")
# 缺域/无事件兜底: 输出列缺失 → null 列(优雅降级,select 不炸)
absent = [c for c in out_cols if c not in out.columns]
if absent:
out = out.with_columns(
[pl.lit(None, dtype=pl.Float64).alias(c) for c in absent])
yield out.sort(["vt_symbol", "datetime"]).select(
["vt_symbol", "datetime", *out_cols])
grid = out = None # 批间释放(下一批重绑定)
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")