Files

160 lines
6.7 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- 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)