e330e130d4
- 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>
1202 lines
59 KiB
Python
1202 lines
59 KiB
Python
# 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_{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")
|
||
|