fix(data): 因子线 review 必修#1 null 桥接洞——首文件全空列把 num/text 桥接(两具体类型从未直接对撞,concat 静默混装零 raise=契约④要防的静默 coerce 被绕过,因子线实测复现);修=每列 resolved=首个非空具体类型/null 永不推进 resolved/冲突即 raise/列集按名比对(列序漂移不误报,策略线观察 b 顺势解决);+3 测试(三段桥 raise/null 隔离仍 raise/列序 tolerated)419 绿;记录②noticed_by 滤空 NOTICE_DATE=有意宁缺毋假写进 docstring;§20.11.5 实施版双线 review 回执+§20.11.6 遗留更新(DSM 复用评估定稿) [vps]
This commit is contained in:
@@ -1427,10 +1427,13 @@ NAS 载荷 corpus 173,003,124+dmsk 220,214,578,另 +2 个 `_mirror_watermark.j
|
||||
0917=5,077 行 / fulltext 按 ids 3/3 命中 / dmsk balance dt=2026-06-30 含
|
||||
noticed_by 过滤 5,572 行、NOTICE_DATE max=披露当日)。
|
||||
|
||||
**遗留**:①DSM 任务 `sanguo-vps-mirror` 每日 16:00 待建(用户 UI 一步:计划
|
||||
任务→用户定义脚本,用户=admin,命令 `bash /volume1/stock/vps-mirror/
|
||||
push_vps_mirror_daily.sh`)——建好前每日增量手动/暂缺;②VPS 侧读取冒烟待
|
||||
vps-deploy 后(读取器代码随 promote 到 VPS);③CLAUDE.md 注记随下次触达补。
|
||||
**遗留**:①每日调度待接——用户倾向**复用现有 DSM 任务**(评估定稿=挂 corpus
|
||||
任务尾部:容器完赛 ~10:1x 即推,同日新鲜度优于独立 16:00;dmsk 07:15 早已
|
||||
完成;失败域用「推送 rc 只记日志不进 corpus 聚合」隔离;03:00 链/dmsk 尾部
|
||||
=方向/时段/新鲜度全错不可选),待用户确认后改 `run_corpus_standalone.sh`
|
||||
尾部+scp 部署;②VPS 侧读取冒烟待 vps-deploy 后(读取器代码随 promote 到
|
||||
VPS);③CLAUDE.md 注记随下次触达补;④backlog:fulltext(ids) 大载荷时加
|
||||
分区定位;稳定 id 内容 hash 若现数据无则采集侧补(另一工序)。
|
||||
|
||||
#### 20.11.5 消费方征询回执(09-17,讨论制:按合理性采用)
|
||||
|
||||
@@ -1462,6 +1465,19 @@ vps-deploy 后(读取器代码随 promote 到 VPS);③CLAUDE.md 注记随
|
||||
(选 fail loud 并入失败分级;实施精化=全空列通配/数值族互通,见 20.11.6
|
||||
④)。镜像链本身无异议。
|
||||
|
||||
**实施版 review 回执(09-18,双线闭环)**:策略线**通过**(三点工程要求逐项
|
||||
核验落地;两条非阻塞观察入 backlog:get_news_fulltext(ids) 全分区扫描随语料
|
||||
增长线性变贵→量大时加 start/end 或按 meta 定位分区;列序漂移误报→已顺势改
|
||||
按列名比对不再误报)。因子线**精化认+必修 1+记录 2**:契约④精化认可(按原
|
||||
字面执行会拒掉全部合法数据);**必修#1 null 桥接洞**(因子线实测复现:首文件
|
||||
全空列把 num 与 text 桥接,两具体类型从未直接对撞→concat 静默混装零 raise,
|
||||
恰是契约④要防的静默 coerce 被绕过)→已修:每列维护 resolved=首个非空具体
|
||||
类型、null 永不推进 resolved、后续非空类型冲突即 raise、列集按名比对(列序
|
||||
漂移不误报),+3 测试(三段桥 raise/null 隔离后续冲突仍 raise/列序漂移
|
||||
tolerated)共 419 绿;记录②NOTICE_DATE 空值行被 noticed_by 滤除=有意
|
||||
宁缺毋假语义(docstring 已写明勿当 bug 修);记录③fulltext 扫描成本同策略
|
||||
线观察,入 backlog。
|
||||
|
||||
## 参考(调查来源)
|
||||
|
||||
- xtdata 官方:https://dict.thinktrader.net/nativeApi/xtdata.html
|
||||
|
||||
@@ -11,6 +11,10 @@ schema 契约(2026-09-17 NAS 实测落档,§20.11.3 数据契约):
|
||||
noticed_by=纯 filter 非 PIT merge(文档明示):只做 NOTICE_DATE<=noticed_by
|
||||
的行级过滤(日期前缀比较);「static 基线为主+dmsk 新鲜度补丁」的 PIT 合并
|
||||
归因子侧 fundamental_pit(§19.3 下游读法),不在数据线。
|
||||
|
||||
noticed_by 语义备注(因子线 09-18 review 记录②):NOTICE_DATE 为空/NaT 的行
|
||||
经前缀比较会被滤除('NaT' 字典序>任何日期串)——**有意语义**=宁缺毋假
|
||||
(无法证明披露可见即不可见),非 bug 勿修。
|
||||
"""
|
||||
import sanguo_data.partition_reader as pr
|
||||
|
||||
|
||||
@@ -13,12 +13,16 @@
|
||||
- 数据在但形状错(列集变化/text↔num 真冲突)→ PartitionSchemaError
|
||||
fail loud(静默 coerce 会埋前视/单位错,研究侧静默错最贵)。
|
||||
|
||||
类型兼容规则(09-17 NAS 真数据实锤精化):EM 源对全空列写 arrow null 型,
|
||||
同域 part 文件间 null↔double/null↔string 漂移是**合法**的——**全空列=通配**
|
||||
(与任意类型兼容);数值族内 int↔float 互通(金额元值 <2^53 无精度损,
|
||||
concat 自然升位);仅 text↔num 冲突与列集增删判漂移。全空 object 列读取时
|
||||
规整为 float64(null↔double 合并后为干净 float64;null↔string 合并为
|
||||
object+NaN 标准缺失形态)。
|
||||
类型兼容规则(09-17 NAS 真数据实锤精化 + 09-18 因子线 review 必修补丁):
|
||||
EM 源对全空列写 arrow null 型,同域 part 文件间 null↔double/null↔string 漂移
|
||||
是**合法**的——**全空列=通配**;数值族内 int↔float 互通(金额元值 <2^53 无
|
||||
精度损,concat 自然升位);仅 text↔num 冲突与列集增删判漂移。**null 桥接洞
|
||||
(review 必修#1)**:通配不能只对首文件两两放行——null 若把 num 与 text 桥接
|
||||
(file1 全空/file2 num/file3 text),concat 会静默混装;修法=每列维护
|
||||
resolved=首个非空具体类型,后续文件非空类型与之冲突即 raise(null 永不推进
|
||||
resolved),列集按名比对(列序漂移不误报)。全空 object 列读取时规整为
|
||||
float64(null↔double 合并后为干净 float64;null↔string 合并为 object+NaN
|
||||
标准缺失形态)。
|
||||
"""
|
||||
import glob
|
||||
import json
|
||||
@@ -79,17 +83,6 @@ def _col_class(series):
|
||||
return str(series.dtype)
|
||||
|
||||
|
||||
def _sig_ok(sig_a, sig_b):
|
||||
if len(sig_a) != len(sig_b):
|
||||
return False
|
||||
for (ca, ta), (cb, tb) in zip(sig_a, sig_b):
|
||||
if ca != cb:
|
||||
return False
|
||||
if ta != tb and ta != "null" and tb != "null":
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def _normalize_null_cols(df):
|
||||
# 全空 object 列规整为 float64:null↔double 合并成干净 float64
|
||||
for col in df.columns:
|
||||
@@ -113,18 +106,28 @@ def read_partitions(root, domain, start=None, end=None):
|
||||
if end is not None:
|
||||
parts = [(d, p) for d, p in parts if d <= end]
|
||||
frames = []
|
||||
sig = None
|
||||
first_names = None
|
||||
resolved = {} # 列名 -> 首个非空具体类型(null 永不推进,防桥接)
|
||||
for _, pdir in parts:
|
||||
for pf in sorted(glob.glob(os.path.join(pdir, "*.parquet"))):
|
||||
df = _normalize_null_cols(pd.read_parquet(pf))
|
||||
cur_sig = tuple((c, _col_class(df[c])) for c in df.columns)
|
||||
if sig is None:
|
||||
sig = cur_sig
|
||||
elif not _sig_ok(sig, cur_sig):
|
||||
bad = [(c, t) for (c, t), (_, t2) in zip(sig, cur_sig)
|
||||
if t != t2 and t != "null" and t2 != "null"] or "列集"
|
||||
cur = [(c, _col_class(df[c])) for c in df.columns]
|
||||
names = set(c for c, _ in cur)
|
||||
if first_names is None:
|
||||
first_names = names
|
||||
if names != first_names:
|
||||
raise PartitionSchemaError(
|
||||
"schema 漂移: %s 冲突列 %s (域=%s)" % (pf, bad, domain))
|
||||
"schema 漂移: %s 列集 %s != 首文件 %s (域=%s)" % (
|
||||
pf, sorted(names), sorted(first_names), domain))
|
||||
for c, t in cur:
|
||||
if t == "null":
|
||||
continue
|
||||
if c not in resolved:
|
||||
resolved[c] = t
|
||||
elif resolved[c] != t:
|
||||
raise PartitionSchemaError(
|
||||
"schema 漂移: %s 列 %s 类型 %s != 已定 %s (域=%s)" % (
|
||||
pf, c, t, resolved[c], domain))
|
||||
frames.append(df)
|
||||
if not frames:
|
||||
return pd.DataFrame()
|
||||
|
||||
@@ -118,6 +118,43 @@ class TestReadPartitions:
|
||||
df = pr.read_partitions(str(root), "x")
|
||||
assert df["v"].dtype == "float64" and sorted(df["v"]) == [2.5, 3.0]
|
||||
|
||||
def test_schema_null_bridging_fails_loud(self, tmp_path):
|
||||
# 因子线 review 必修#1:首文件全空列把 num 与 text 桥接,三段必须 raise
|
||||
root = tmp_path / "d"
|
||||
(root / "x" / "dt=2026-06-30").mkdir(parents=True)
|
||||
pd.DataFrame({"v": [None], "k": ["a"]}).to_parquet(
|
||||
root / "x" / "dt=2026-06-30" / "part-0.parquet", index=False)
|
||||
pd.DataFrame({"v": [1], "k": ["b"]}).to_parquet(
|
||||
root / "x" / "dt=2026-06-30" / "part-1.parquet", index=False)
|
||||
pd.DataFrame({"v": ["txt"], "k": ["c"]}).to_parquet(
|
||||
root / "x" / "dt=2026-06-30" / "part-2.parquet", index=False)
|
||||
with pytest.raises(pr.PartitionSchemaError):
|
||||
pr.read_partitions(str(root), "x")
|
||||
|
||||
def test_schema_null_between_conflicts_still_fails(self, tmp_path):
|
||||
# resolved 不被 null 文件重置:num→null→text 同样 raise
|
||||
root = tmp_path / "d"
|
||||
(root / "x" / "dt=2026-06-30").mkdir(parents=True)
|
||||
pd.DataFrame({"v": [1], "k": ["a"]}).to_parquet(
|
||||
root / "x" / "dt=2026-06-30" / "part-0.parquet", index=False)
|
||||
pd.DataFrame({"v": [None], "k": ["b"]}).to_parquet(
|
||||
root / "x" / "dt=2026-06-30" / "part-1.parquet", index=False)
|
||||
pd.DataFrame({"v": ["txt"], "k": ["c"]}).to_parquet(
|
||||
root / "x" / "dt=2026-06-30" / "part-2.parquet", index=False)
|
||||
with pytest.raises(pr.PartitionSchemaError):
|
||||
pr.read_partitions(str(root), "x")
|
||||
|
||||
def test_schema_column_order_drift_tolerated(self, tmp_path):
|
||||
# 同列集不同列序:按列名对齐不误报(策略线 review 观察顺势处理)
|
||||
root = tmp_path / "d"
|
||||
(root / "x" / "dt=2026-06-30").mkdir(parents=True)
|
||||
pd.DataFrame({"a": [1], "b": ["x"]}).to_parquet(
|
||||
root / "x" / "dt=2026-06-30" / "p0.parquet", index=False)
|
||||
pd.DataFrame({"b": ["y"], "a": [2]}).to_parquet(
|
||||
root / "x" / "dt=2026-06-30" / "p1.parquet", index=False)
|
||||
df = pr.read_partitions(str(root), "x")
|
||||
assert len(df) == 2 and df["a"].tolist() == [1, 2]
|
||||
|
||||
|
||||
class TestWatermark:
|
||||
def test_read_watermark_dict(self, tree):
|
||||
|
||||
Reference in New Issue
Block a user