# -*- coding: utf-8 -*- """corpus 语料薄封装(spec §20.11.3)。列名随源直存零转换(§15.2 双层读法)。 schema 契约(2026-09-17 NAS 实测落档,§20.11.3 数据契约③): - news_meta: art_code(稳定 id)/stock_code(裸 6 位)/show_time/title/summary/ media_name/url/first_seen_date - news_fulltext: art_code(=news_meta join 键)/url/show_time/content_text/fetch_date - ann_meta: announcement_id(稳定 id,与 ann_pdf//.pdf 同名)/ sec_code(裸 6 位)/sec_name/title/short_title/ann_time/ann_type(|| 分隔多码)/ ann_type_name/column_id/adjunct_url/adjunct_size/first_seen_date - ann_pdf: /ann_pdf//.pdf - flash_meta(§20.12.3, 2026-09-22): flash_id(jin10:/em_kuaixun: 前缀+源站稳定 id)/ source/publish_time/title/content/important/channel/first_seen_date—— 全市场宏观快讯流(金十+东财kuaixun), 与 news_meta 个股新闻语义并列不 join 稳定 id(数据契约②配方):art_code/announcement_id 为源站稳定 id;爬虫层 跨天重复行会在数据里原样出现——核内零去重,消费方按 id 自行幂等。 """ import os import sanguo_data.partition_reader as pr def get_news_meta(start=None, end=None, symbols=None): """新闻元数据(dt 分区直读)。symbols=裸 6 位代码列表,下推过滤 stock_code。""" df = pr.read_partitions(pr.corpus_root(), "news_meta", start, end) if symbols and not df.empty: df = df[df["stock_code"].isin(set(symbols))] return df def get_ann_meta(start=None, end=None, symbols=None): """公告元数据。symbols 下推过滤 sec_code;ann_type/标题已透传(因子线 选择性抽取 PDF 用,不盲抽)。""" df = pr.read_partitions(pr.corpus_root(), "ann_meta", start, end) if symbols and not df.empty: df = df[df["sec_code"].isin(set(symbols))] return df def get_news_fulltext(ids): """按 art_code 批量取全文(因子回访高信号条目/LLM 分块天然按 id)。""" df = pr.read_partitions(pr.corpus_root(), "news_fulltext") if ids and not df.empty: df = df[df["art_code"].isin(set(ids))] return df def ann_pdf_dir(year): """公告 PDF 年目录(文件名=announcement_id.pdf,与 ann_meta 同名 join)。""" return os.path.join(pr.corpus_root(), "ann_pdf", str(year)) FLASH_COLUMNS = ["flash_id", "source", "publish_time", "title", "content", "important", "channel", "first_seen_date"] def get_flash(start=None, end=None, sources=None): """全市场快讯流(§20.12.3,dt 分区直读)。sources={"jin10","em_kuaixun"} 过滤;空树返带 schema 的空 df(容错分级契约)。""" df = pr.read_partitions(pr.corpus_root(), "flash_meta", start, end) if df.empty: return df.reindex(columns=FLASH_COLUMNS) if sources: df = df[df["source"].isin(set(sources))] return df EVENTS_COLUMNS = ["event_id", "content_hash", "src_domain", "src_id", "stock_code", "event_date", "event_type", "direction", "confidence", "model", "extracted_at"] def get_events(start=None, end=None, models=None): """事件判读快照(spec §20.16 P4-3,dt 分区直读)。models 过滤 ({"lexicon"}=只看词表层,或模型名);空树返带 schema 的空 df (容错分级契约同 get_flash)。""" df = pr.read_partitions(pr.corpus_root(), "events_llm", start, end) if df.empty: return df.reindex(columns=EVENTS_COLUMNS) if models: df = df[df["model"].isin(set(models))] return df def get_watermark(): """镜像水位线(None=尚无成功推送记录)。""" return pr.read_watermark(pr.corpus_root())