Files
sanguo_vnpy_v2/sanguo_factor/fundamental_adapter.py
T

814 lines
40 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 自算(表达式层)
输出: vt_symbol × datetime(日频) × FEATURE_COLUMNS,join 到 alpha_df 后
供 cs_rank(列) 表达式直接消费。
"""
from __future__ import annotations
import os
import warnings
from datetime import 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"]
# 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",
]
# 金融股置 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"]
# forecast 事件流日频列(公告日 asof): P0 两列 + P1 年化预告净利(D13)
_FORECAST_COLS = ("forecast_type_score", "forecast_change_pct",
"forecast_np_annualized")
# 预告兑现差独立事件流(F07 前视红线: 锚 = max(实际披露日, 预告公告日))
_BEAT_COLS = ("forecast_beat",)
# ==================== 读取层 ====================
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"),
)
# 第三级: 行本地比率(无时序,无需 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"),
)
# 第四级: 跨期差分/同比/剪刀差(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(
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
# ==================== 对外主入口 ====================
# 分块大小: 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]
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)
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")
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")