domain 模块

August 29, 2026 · View on GitHub

路径:src/cnequity/domain/

数据契约层:schema 类型、主键、数据集元数据、符号规则、跨进程限速、情绪打分工具。不含 I/O 与编排。


文件一览

文件职责
schemas.pyPolars schema、PRIMARY_KEYSvalidate_dataframe()with_provenance()
datasets.pyDatasetSpec 注册表 DATASETS 与数据集元数据
contracts.py稳定 JSON 契约、fingerprint、验证和 breaking diff
symbols.pyparse_symbol(), is_all_a_symbol(), CDR/ETF 分类
rate_limit.py跨平台文件锁 + JSON 时隙状态的跨进程 RateLimiter
sentiment.py公告关键词 + 可选 SnowNLP 打分

schemas.py

核心常量

  • DATASET_SCHEMAS: dict[str, dict[str, pl.DataType]] — 每数据集列类型
  • PRIMARY_KEYS: dict[str, list[str]] — 主键列
  • MOCK_SOURCE = "mock" — 测试源标识

溯源列

每个 curated schema 末尾包含:

"source": pl.Utf8,
"data_version": pl.Utf8,
"fetched_at": pl.Datetime("us", "UTC"),

validate_dataframe(df, dataset)

写 staging/curated 前调用:

  • 列齐全且类型匹配
  • 不允许未知列(strict)
  • PK、sourcedata_versionfetched_at 必须非空;字符串溯源/主键不能是空白
  • 所有浮点字段拒绝 NaN / Inf;行情数据额外校验价格、成交量与 OHLC 关系

with_provenance(df, source, data_version)

为 adapter 输出批量添加溯源列。


datasets.py

DatasetSpec 字段

字段含义
name数据集名
tierL0–L8 研究分层;无默认值,必须显式声明
layercurated / derived存储位置,与 tier 正交)
partition_colHive 分区列;None = merge 文件
partition_granularityday / month / quarter / year;按每日行数选,不按习惯
date_col查询日期列;默认等于 partition_col
fetch_semanticsby_date / snapshot
watermark是否维护 meta/state 水位
pit是否 PIT 数据集
backfill_sourcesnapshot 数据集的历史回填源名
max_staleness_daysstatus --datasets 容忍滞后天数
requiredFalse 时空 curated 只算 warning,不拉低 lake_health
history_horizon_days源端还提供多少个交易日(滚动,随今天前移)
history_floor_date源端的固定日历底(不随今天移动);与上一项二选一,同时设时它优先
backfill_chunk_days单次回填子跑覆盖的日历天数(by-date 源用)
backfill_chunk_symbols单次回填子跑的标的数(tip-paged 源用,与上一项互斥)
intraday_frequencybar 频率(1m / 5m)。行为字段:设了就会被 audit 的会话检查、reader 的复权集合、cne backfill --symbols 认领
row_grain一行覆盖多久(1m / 5m / tick)。纯描述,不驱动任何行为
coverage_modesession_dense 表示覆盖区间内每个交易日都应有数据;sparse 表示事件/公告等允许空交易日
schema_version列形状版本;默认 1,新增列保持向后兼容
contract_level契约成熟度;当前注册表默认为 stable
pit_grade0.x 兼容别名:none / strict / partial;历史回填为 partial
pit_qualitystrict / reconstructed / snapshot_only;描述当前证据质量
availability_col信息可用日期列;PIT 默认 announce_date,其他表默认查询日期列
pit_storage_columns可选双时态列:available_atsource_published_atobserved_atrevision_id
compatibility兼容策略;现有表默认 additive
unit_contract数值字段的机器可读单位声明

两组容易混的字段:

history_horizon_days vs history_floor_date —— 前者是「每标的固定根数」(分钟线:源端存 22,800 根 1m,除以一个完整交易日得 95 天),窗口每天往前滑;后者是服务端按日历切的保留底(分笔:所有标的都回溯到 2024-01-02),不随今天移动,所以视野逐日变长。用错会让 earliest_available() 每天漂,几个月后把源端还愿意给的数据挡在门外。

intraday_frequency vs row_grain —— 前者是行为的,它的消费者都假定存在 bar_time 列和「每交易日 N 根」;trade_ticks 故意不设它,否则会继承一批在错误列上静默通过的检查。但目录和面板仍需知道它是日内数据,这是 row_grain 的唯一职责。两者同时存在时必须一致(注册表测试强制)。

辅助函数

get_dataset(name) -> DatasetSpec
curated_dataset_names() -> frozenset[str]
derived_dataset_names() -> frozenset[str]
pit_dataset_names() -> frozenset[str]
fetch_semantics(dataset) -> Literal["by_date", "snapshot"]
is_stale(dataset, mark, anchor) -> bool

完整 JSON 契约及演进比较见 datasets/contract.md。 程序化入口为 dataset_contract()build_contract()contract_fingerprint()validate_contract()diff_contracts()

新增数据集必须:在此添加 DatasetSpec + 在 schemas.py 添加 schema/PK。tests/unit/test_dataset_registry.py 强制同步。


symbols.py

  • 格式:600519.SH000001.SZ920001.BJ
  • is_all_a_symbol():沪深 A 股前缀白名单(60/68/00/30/92);不含 ETF 前缀
  • is_cdr_symbol():SH 689 段存托凭证
  • is_etf_symbol():场内 ETF/LOF(SH 51/52/56/58,SZ 15/16
  • instruments 抓取含股票 + CDR + ETF;parse_symbol()(code, exchange)

Universe 过滤在 query/universe.py 使用本模块规则:all_a 排除 CDR 与 ETF。


rate_limit.py

跨进程限速:共享 JSON 状态保存 next_allowed_at,进程先在短锁事务内预订自己的请求时隙,再释放锁并在锁外等待。这样多个 worker 仍共享同一源级请求间隔,但一个 worker 的 sleep 不会把其他 worker 堵在文件锁上。旧的仅含 last 的状态文件会自动迁移;锁等待超过 lock_timeout_sec 会显式失败,不会绕过限速。锁通过包根 file_lock.exclusive_lock(POSIX flock / Windows msvcrt.locking)取得。


sentiment.py

research step 使用:

  • 公告标题/正文关键词情绪
  • use_snownlp=true 时调用 SnowNLP(可选依赖)

相关文档