domain 模块
August 29, 2026 · View on GitHub
路径:src/cnequity/domain/
数据契约层:schema 类型、主键、数据集元数据、符号规则、跨进程限速、情绪打分工具。不含 I/O 与编排。
文件一览
| 文件 | 职责 |
|---|---|
schemas.py | Polars schema、PRIMARY_KEYS、validate_dataframe()、with_provenance() |
datasets.py | DatasetSpec 注册表 DATASETS 与数据集元数据 |
contracts.py | 稳定 JSON 契约、fingerprint、验证和 breaking diff |
symbols.py | parse_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、
source、data_version、fetched_at必须非空;字符串溯源/主键不能是空白 - 所有浮点字段拒绝
NaN/Inf;行情数据额外校验价格、成交量与 OHLC 关系
with_provenance(df, source, data_version)
为 adapter 输出批量添加溯源列。
datasets.py
DatasetSpec 字段
| 字段 | 含义 |
|---|---|
name | 数据集名 |
tier | L0–L8 研究分层;无默认值,必须显式声明 |
layer | curated / derived(存储位置,与 tier 正交) |
partition_col | Hive 分区列;None = merge 文件 |
partition_granularity | day / month / quarter / year;按每日行数选,不按习惯 |
date_col | 查询日期列;默认等于 partition_col |
fetch_semantics | by_date / snapshot |
watermark | 是否维护 meta/state 水位 |
pit | 是否 PIT 数据集 |
backfill_source | snapshot 数据集的历史回填源名 |
max_staleness_days | status --datasets 容忍滞后天数 |
required | False 时空 curated 只算 warning,不拉低 lake_health |
history_horizon_days | 源端还提供多少个交易日(滚动,随今天前移) |
history_floor_date | 源端的固定日历底(不随今天移动);与上一项二选一,同时设时它优先 |
backfill_chunk_days | 单次回填子跑覆盖的日历天数(by-date 源用) |
backfill_chunk_symbols | 单次回填子跑的标的数(tip-paged 源用,与上一项互斥) |
intraday_frequency | bar 频率(1m / 5m)。行为字段:设了就会被 audit 的会话检查、reader 的复权集合、cne backfill --symbols 认领 |
row_grain | 一行覆盖多久(1m / 5m / tick)。纯描述,不驱动任何行为 |
coverage_mode | session_dense 表示覆盖区间内每个交易日都应有数据;sparse 表示事件/公告等允许空交易日 |
schema_version | 列形状版本;默认 1,新增列保持向后兼容 |
contract_level | 契约成熟度;当前注册表默认为 stable |
pit_grade | 0.x 兼容别名:none / strict / partial;历史回填为 partial |
pit_quality | strict / reconstructed / snapshot_only;描述当前证据质量 |
availability_col | 信息可用日期列;PIT 默认 announce_date,其他表默认查询日期列 |
pit_storage_columns | 可选双时态列:available_at、source_published_at、observed_at、revision_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.SH、000001.SZ、920001.BJ is_all_a_symbol():沪深 A 股前缀白名单(60/68/00/30/92);不含 ETF 前缀is_cdr_symbol():SH689段存托凭证is_etf_symbol():场内 ETF/LOF(SH51/52/56/58,SZ15/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(可选依赖)