89 lines
3.7 KiB
Python
89 lines
3.7 KiB
Python
# -*- 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/<year>/<id>.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: <corpus_root>/ann_pdf/<year>/<announcement_id>.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())
|