160 lines
6.7 KiB
Python
160 lines
6.7 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""通用 dt= 分区 parquet 读取核(spec §20.11.3,corpus/dmsk 消费接口基座)。
|
||
|
||
核内零语义(因子线 09-17 契约):
|
||
- 不做 as-of/PIT 合并——因子侧 fundamental_pit 是唯一正确性模型,两处实现会漂
|
||
且 PIT 错是静默错;dmsk 按 REPORT_DATE 分区,「as of D」本就无法靠分区剪枝。
|
||
- 零静默去重——幂等去重归写路径;镜像重放/分区重写出重复行原样透出,
|
||
取哪行=消费方语义决策。
|
||
|
||
失败分级(策略线 09-17 契约):
|
||
- 缺数据(根/域/分区不存在)→ 空 DataFrame,绝不 raise
|
||
(「数据失败≠交易信号」,消费方沿用上次快照);
|
||
- 分区残骸(0 字节=掉电写盘中断遗物)→ 跳过该 part+ERROR 日志
|
||
不裸崩(10-05 NAS 掉电案 funnel walk 实锤;写侧 fsync 已根治,此为纵深;
|
||
**坏件含内容者继续 raise**——integrity_gate 单查隔离依赖异常记 error=);
|
||
- 数据在但形状错(列集变化/text↔num 真冲突)→ PartitionSchemaError
|
||
fail loud(静默 coerce 会埋前视/单位错,研究侧静默错最贵)。
|
||
|
||
类型兼容规则(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
|
||
import logging
|
||
import os
|
||
|
||
import pandas as pd
|
||
|
||
log = logging.getLogger("partition_reader")
|
||
|
||
_IS_WIN = os.name == "nt"
|
||
DEFAULT_CORPUS_ROOT_WIN = r"C:\sanguo_vnpy_v2\data\corpus"
|
||
DEFAULT_CORPUS_ROOT_POSIX = "/volume1/stock/corpus"
|
||
DEFAULT_DMSK_ROOT_WIN = r"C:\sanguo_vnpy_v2\data\fundamentals\dmsk"
|
||
DEFAULT_DMSK_ROOT_POSIX = "/volume1/stock/fundamentals/dmsk"
|
||
WATERMARK_NAME = "_mirror_watermark.json"
|
||
_DT_PREFIX = "dt="
|
||
|
||
|
||
class PartitionSchemaError(RuntimeError):
|
||
"""跨分区 schema 漂移(列集/类型不一致)。"""
|
||
|
||
|
||
def corpus_root():
|
||
"""corpus 族根(§20.11.3:双机各读各的,env 覆盖同 SANGUO_DATA_ROOT 模式)。"""
|
||
return os.environ.get("SANGUO_CORPUS_ROOT") or (
|
||
DEFAULT_CORPUS_ROOT_WIN if _IS_WIN else DEFAULT_CORPUS_ROOT_POSIX)
|
||
|
||
|
||
def fund_dmsk_root():
|
||
"""dmsk 族根。"""
|
||
return os.environ.get("SANGUO_FUND_DMSK_ROOT") or (
|
||
DEFAULT_DMSK_ROOT_WIN if _IS_WIN else DEFAULT_DMSK_ROOT_POSIX)
|
||
|
||
|
||
def _partitions(domain_root):
|
||
if not os.path.isdir(domain_root):
|
||
return []
|
||
out = []
|
||
for name in sorted(os.listdir(domain_root)):
|
||
p = os.path.join(domain_root, name)
|
||
if os.path.isdir(p) and name.startswith(_DT_PREFIX):
|
||
out.append((name[len(_DT_PREFIX):], p))
|
||
return out
|
||
|
||
|
||
def list_dates(root, domain):
|
||
"""可用分区日期列表(dt= 前缀剥掉,升序)。缺根/缺域=空列表。"""
|
||
return [d for d, _ in _partitions(os.path.join(root, domain))]
|
||
|
||
|
||
def _col_class(series):
|
||
"""列类型兼容类:全空列=通配 null;数值族互通;text/bool 精确。"""
|
||
if series.isna().all():
|
||
return "null"
|
||
k = series.dtype.kind
|
||
if k in "ifu":
|
||
return "num"
|
||
if k in "OSU":
|
||
return "text"
|
||
return str(series.dtype)
|
||
|
||
|
||
def _normalize_null_cols(df):
|
||
# 全空 object 列规整为 float64:null↔double 合并成干净 float64
|
||
for col in df.columns:
|
||
s = df[col]
|
||
if s.dtype == object and s.isna().all():
|
||
df[col] = s.astype("float64")
|
||
return df
|
||
|
||
|
||
def read_partitions(root, domain, start=None, end=None):
|
||
"""读域的 dt= 分区并拼接(含分区内全部 part-*.parquet)。
|
||
|
||
start/end 为 ISO 日期字符串,含端点、字符串比较。零去重零转换;
|
||
列集增删或 text↔num 冲突 raise PartitionSchemaError(全空列通配/
|
||
数值族内 int↔float 互通=合法 EM null 型漂移,见模块 docstring);
|
||
无数据返空 DataFrame。
|
||
"""
|
||
parts = _partitions(os.path.join(root, domain))
|
||
if start is not None:
|
||
parts = [(d, p) for d, p in parts if d >= start]
|
||
if end is not None:
|
||
parts = [(d, p) for d, p in parts if d <= end]
|
||
frames = []
|
||
first_names = None
|
||
resolved = {} # 列名 -> 首个非空具体类型(null 永不推进,防桥接)
|
||
for _, pdir in parts:
|
||
for pf in sorted(glob.glob(os.path.join(pdir, "*.parquet"))):
|
||
# 10-05 掉电防御(读侧): 仅 0 字节残骸(写盘中断遗物, 写侧 fsync
|
||
# 已根治此为纵深)跳过+ERROR——坏件含内容者继续 fail loud:
|
||
# integrity_gate 的单查隔离设计依赖捕获异常记 error=(1ea143e
|
||
# 实测打穿教训, 勿再扩 try 面)
|
||
if os.path.getsize(pf) == 0:
|
||
log.error("分区残骸跳过(0 字节, 掉电/写盘中断遗物): %s", pf)
|
||
continue
|
||
df = _normalize_null_cols(pd.read_parquet(pf))
|
||
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 (域=%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()
|
||
return pd.concat(frames, ignore_index=True)
|
||
|
||
|
||
def read_watermark(root):
|
||
"""镜像水位线(§20.11.2 推送链落在族根的 _mirror_watermark.json)。
|
||
|
||
返回 {generated_at, domains:{domain: 最大分区值}} 或 None(无文件)。
|
||
"""
|
||
p = os.path.join(root, WATERMARK_NAME)
|
||
if not os.path.exists(p):
|
||
return None
|
||
with open(p) as f:
|
||
return json.load(f)
|