feat(data): 漏斗主体 corpus_funnel——两层分诊(词表走水位+LLM lane 串行预算断路)+事件快照两域+pending 持久队列+卡片需求面; Ctx 加 parse_fail 计数, ID_KEYS 补 events 两域(P4-3 F2/F3/F5, 计划 Task4+5 合落=模块内聚) [nas] [no-doc]

This commit is contained in:
2026-10-02 09:54:21 +08:00
parent 8bf935b1a0
commit 65e5c02916
3 changed files with 619 additions and 1 deletions
+5 -1
View File
@@ -89,7 +89,10 @@ UA = {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.3
ID_KEYS = {"ann_meta": ["announcement_id"],
"news_meta": ["art_code", "stock_code"],
"news_fulltext": ["art_code"],
"flash_meta": ["flash_id"]}
"flash_meta": ["flash_id"],
# P4-3 事件快照两域(spec §4.8 F2): 输入切片按内容寻址, 判读按事件 id
"events_in": ["content_hash"],
"events_llm": ["event_id"]}
FLASH_EM_MAX_PAGES = 20 # 快照型时效数据: 50条/页×20≈2-3天窗, 不深回补
FLASH_LOOKBACK_DAYS = 1 # 东财快讯翻页拉到昨日 00:00 边界即停
@@ -736,6 +739,7 @@ class Ctx:
self.deadline = dt.datetime.combine(
dt.date.today(), dt.time(int(hh), int(mm)))
self.failed = 0
self.parse_fail = 0 # 解析/契约违规计数(P4-3 漏斗)
self.rate_limited = False
self.hard_cool = False
self.wallclock = False
+375
View File
@@ -0,0 +1,375 @@
# -*- coding: utf-8 -*-
"""corpus_funnel.py — 漏斗分诊+事件快照夜批(spec §4.8 决议 G/C7+F1-F5, P4-3).
两层: 第一层确定性词表(funnel_lexicon, 零成本处理大头) → 词表 miss 且含
信号词 → 第二层 LLM(只碰漏下来的, sanguo_api/llm client 串行慢爬). 产物=
事件快照两域(append-only dt 分区=采集日, 永不重算):
events_in LLM 候选输入切片(防编数: 喂模型的确切文本先落盘, 可回放对拍)
events_llm 事件判读(词表+LLM 统一域, model 列区分来源)
基建全复用 corpus_download(Ctx 预算断路/_absorb_new 账本去重/marker);
自有 flock(state/funnel.lock, corpus 班错峰后跑互不锁). 增量=state/
funnel_watermark.json 各源域独立水位只前向; 历史回填=--start 按预算分批;
LLM 候选持久队列 state/funnel_pending.jsonl(预算停摆残余次夜续).
退出码(与 corpus_download 同契约): 0=完成 1=致命(unit 失败/LLM 配置缺失
且有候选) 2=限流让路 3=预算/墙钟让路(checkpoint, 次夜续).
用法: python corpus_funnel.py --until 23:00 --llm-cap 200
python corpus_funnel.py --no-llm # 只跑第一层
python corpus_funnel.py --start 2026-09-06 # 历史回填(按预算分批)
"""
from __future__ import annotations
import argparse
import asyncio
import datetime as dt
import fcntl
import hashlib
import json
import logging
import os
import sys
from pathlib import Path
_HERE = os.path.dirname(os.path.abspath(__file__))
_REPO = os.path.dirname(os.path.dirname(_HERE))
for _p in (_HERE, _REPO):
if _p not in sys.path:
sys.path.insert(0, _p)
import pandas as pd # noqa: E402
from scripts.data_platform import corpus_download as cd # noqa: E402
from scripts.data_platform import funnel_lexicon as lex # noqa: E402
from sanguo_data import partition_reader as pr # noqa: E402
log = logging.getLogger("funnel")
PENDING_NAME = "funnel_pending.jsonl"
WATERMARK_NAME = "funnel_watermark.json"
SRC_DOMAINS = ("news_meta", "ann_meta")
LLM_CONSEC_FAIL_BREAK = 3 # 连续失败熔断(防烧夜, 已处理件有 marker 不重)
EVENTS_IN_COLUMNS = ["content_hash", "src_domain", "src_id", "stock_code",
"event_date", "title", "text_snippet", "selected_at"]
EVENTS_LLM_COLUMNS = ["event_id", "content_hash", "src_domain", "src_id",
"stock_code", "event_date", "event_type", "direction",
"confidence", "model", "extracted_at"]
PROMPT_SYSTEM = (
"你是A股公告/新闻事件抽取器。输入是单一股票单一日期的一段标题与摘要。"
"只依据输入文本判断, 禁止编造输入中没有的信息。"
'输出严格 JSON: {"events": [{"event_type": "...", '
'"direction": "pos|neg|neutral", "confidence": 0.0到1.0}]}。'
"event_type 只能取以下之一: "
+ ", ".join(f"{k}({v})" for k, v in lex.EVENT_TYPES.items())
+ "。无把握或无事件时输出空 events 数组。stock_code 与 event_date 由系统"
"从输入行带入, 你不要输出它们。")
# ---------- 快照行构造 ----------
def _content_hash(domain, src_id, code, title):
s = f"{domain}|{src_id}|{code}|{title}"
return hashlib.blake2b(s.encode("utf-8"), digest_size=8).hexdigest()
def _event_row(content_hash, src_domain, src_id, code, edate, hit, model):
eid = hashlib.blake2b(
f"{content_hash}|{hit['event_type']}|{code}|{hit['direction']}"
.encode("utf-8"), digest_size=8).hexdigest()
now = dt.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
return {"event_id": eid, "content_hash": content_hash,
"src_domain": src_domain, "src_id": str(src_id),
"stock_code": str(code), "event_date": edate,
"event_type": hit["event_type"], "direction": hit["direction"],
"confidence": float(hit["confidence"]), "model": model,
"extracted_at": now}
def _extract_day(domain, day, df):
"""一天 corpus 行 → (词表事件行, LLM 候选输入行)。"""
events, cands = [], []
for row in df.to_dict("records"):
if domain == "news_meta":
src_id, code = row.get("art_code"), row.get("stock_code")
when, title = row.get("show_time"), row.get("title")
text = row.get("summary")
else:
src_id, code = row.get("announcement_id"), row.get("sec_code")
when, title = row.get("ann_time"), row.get("title")
text = row.get("content") or row.get("short_title")
if not src_id or not code or not title:
continue
edate = str(when or "")[:10] or day
ch = _content_hash(domain, src_id, code, title)
for hit in lex.match_events(title, text):
events.append(_event_row(ch, domain, src_id, code, edate, hit,
model="lexicon"))
if lex.is_llm_candidate(title, text):
cands.append({"content_hash": ch, "src_domain": domain,
"src_id": str(src_id), "stock_code": str(code),
"event_date": edate, "title": title,
"text_snippet": f"{title}\n{text or ''}"[:500],
"selected_at": dt.date.today().isoformat()})
return events, cands
# ---------- 水位与 pending 队列 ----------
def _load_watermark():
p = cd.CORPUS_ROOT / "state" / WATERMARK_NAME
if p.exists():
try:
return json.loads(p.read_text(encoding="utf-8"))
except ValueError:
log.warning("%s 损坏, 重置", WATERMARK_NAME)
return {}
def _save_watermark(wm):
p = cd.CORPUS_ROOT / "state" / WATERMARK_NAME
p.parent.mkdir(parents=True, exist_ok=True)
tmp = p.with_suffix(".tmp")
tmp.write_text(json.dumps(wm), encoding="utf-8")
os.replace(tmp, p)
def _pending_path():
return cd.CORPUS_ROOT / "state" / PENDING_NAME
def _load_pending():
p = _pending_path()
if not p.exists():
return []
out = []
for ln in p.read_text(encoding="utf-8").splitlines():
ln = ln.strip()
if not ln:
continue
try:
row = json.loads(ln)
except ValueError:
continue
if not cd.is_done("funnel", "llm", row.get("content_hash")):
out.append(row) # 消费 marker 幂等: done 件不重排
return out
def _save_pending(rows):
p = _pending_path()
p.parent.mkdir(parents=True, exist_ok=True)
tmp = p.with_suffix(".tmp")
tmp.write_text(
"".join(json.dumps(r, ensure_ascii=False) + "\n" for r in rows),
encoding="utf-8")
os.replace(tmp, p)
# ---------- 单实例锁(flock, corpus 班错峰互不干扰) ----------
def _acquire_lock():
lock = cd.CORPUS_ROOT / "state" / "funnel.lock"
lock.parent.mkdir(parents=True, exist_ok=True)
fd = os.open(str(lock), os.O_CREAT | os.O_RDWR, 0o666)
try:
fcntl.flock(fd, fcntl.LOCK_EX | fcntl.LOCK_NB)
except BlockingIOError:
os.close(fd)
return None
return fd
# ---------- 卡片需求面(F5: 记录不驱动) ----------
def _read_card_demand(pipeline_db):
if not pipeline_db or not os.path.exists(pipeline_db):
return []
try:
from sanguo_portfolio import pipeline_store
items = pipeline_store.list_hypotheses(pipeline_db)
except Exception as e: # noqa: BLE001 读侧可选件
log.warning("卡片需求面读取失败(非致命): %s", e)
return []
return [{"hyp_id": i.get("id"), "dataNeeds": i.get("dataNeeds")}
for i in items if i.get("state") == "data_check"]
# ---------- 第一层: 词表走水位 ----------
def _walk_domain(ctx, domain, args, ledgers, recents, fresh_cands, stats):
root = cd.CORPUS_ROOT
wm = _load_watermark()
todo = [d for d in pr.list_dates(str(root), domain)
if (wm.get(domain) is None or d > wm[domain])
and (args.start is None or d >= args.start)]
for d in sorted(todo):
ctx.stop_now() # WallClock/BudgetStop
df = pr.read_partitions(str(root), domain, start=d, end=d)
if not df.empty:
events, cands = _extract_day(domain, d, df)
if events:
cd._absorb_new("events_llm", events, ledgers, recents)
stats["lex_events"] += len(events)
fresh_cands.extend(cands[: args.llm_cap])
stats["candidates_fresh"] += len(cands[: args.llm_cap])
ctx.unit_done(domain, d)
stats["days"] += 1
wm[domain] = d
_save_watermark(wm)
# ---------- 第二层: LLM lane(串行慢爬, 预算断路) ----------
async def _llm_lane(ctx, client, model_name, cands, ledgers, recents):
"""逐候选串行抽取。返回写出事件数。
- 已 done(content_hash marker)件跳过(崩溃恢复幂等);
- 预算触顶抛 BudgetStop 由调用方收拾(残余回写 pending);
- 非法判读项确定性丢弃计 ctx.parse_fail, 不整批废(F3);
- 消费即落 events_in(防编数: 喂模型文本先落盘)。
"""
written = 0
consec_fail = 0
for row in cands:
if cd.is_done("funnel", "llm", row["content_hash"]):
continue
ctx.stop_now()
try:
resp, usage = await client.chat_json(
[{"role": "system", "content": PROMPT_SYSTEM},
{"role": "user", "content": _prompt_user(row)}],
want_usage=True)
except Exception as e: # noqa: BLE001 LLMError 族
log.warning("llm fail %s: %s", row["content_hash"], e)
ctx.failed += 1
consec_fail += 1
if consec_fail >= LLM_CONSEC_FAIL_BREAK:
log.error("llm 连续 %d 败, 熔断余量", consec_fail)
break
continue
consec_fail = 0
ctx.add_tokens(usage.get("prompt_tokens"), usage.get("completion_tokens"))
events = []
for ev in (resp.get("events") or []) if isinstance(resp, dict) else []:
etype = ev.get("event_type")
direction = ev.get("direction")
try:
conf = float(ev.get("confidence"))
except (TypeError, ValueError):
conf = -1.0
if (etype not in lex.EVENT_TYPES
or direction not in ("pos", "neg", "neutral")
or not 0.0 <= conf <= 1.0):
ctx.parse_fail += 1
continue
events.append(_event_row(
row["content_hash"], row["src_domain"], row["src_id"],
row["stock_code"], row["event_date"],
{"event_type": etype, "direction": direction,
"confidence": conf}, model_name))
if events:
cd._absorb_new("events_llm", events, ledgers, recents)
written += len(events)
cd._absorb_new("events_in", [dict(row)], ledgers, recents)
ctx.unit_done("llm", row["content_hash"])
return written
def _prompt_user(row):
return (f"stock_code={row['stock_code']}\nevent_date={row['event_date']}\n"
f"标题: {row['title']}\n摘要: {row['text_snippet']}")
def _dump_raw_fail(content_hash, resp_text):
d = cd.CORPUS_ROOT / "state" / "funnel_raw"
d.mkdir(parents=True, exist_ok=True)
(d / f"{content_hash}.txt").write_text(str(resp_text), encoding="utf-8")
# ---------- 入口 ----------
def main(argv=None):
ap = argparse.ArgumentParser(description=__doc__)
ap.add_argument("--until", help="墙钟自停 HH:MM(rc=3 checkpoint)")
ap.add_argument("--limit", type=int, help="LLM 调用次数上限(冒烟不落 marker)")
ap.add_argument("--token-budget", type=int, help="夜批 token 上限(rc=3)")
ap.add_argument("--llm-cap", type=int, default=200,
help="每夜新鲜候选上限(默认 200)")
ap.add_argument("--start", help="回填起点 YYYY-MM-DD(默认=水位续跑)")
ap.add_argument("--no-llm", action="store_true", help="只跑第一层")
ap.add_argument("--pipeline-db", default=os.environ.get("SANGUO_PIPELINE_DB"),
help="卡片库路径(data_check 态需求面, 缺省 env/缺文件=跳过)")
args = ap.parse_args(argv)
logging.basicConfig(level=logging.INFO,
format="%(asctime)s %(levelname)s %(name)s %(message)s")
fd = _acquire_lock()
if fd is None:
print("[funnel] 锁让路退出(另一实例在跑)")
return 3
ctx = cd.Ctx("funnel", until=args.until, limit=args.limit,
token_budget=args.token_budget)
ledgers = {d: cd.IdLedger(d) for d in ("events_in", "events_llm")}
recents = {d: cd.RecentIndex(d) for d in ("events_in", "events_llm")}
ctx.set_stores(ledgers)
cd._reconcile_ledgers(ledgers, recents)
stats = {"days": 0, "lex_events": 0, "candidates_fresh": 0,
"llm_calls": 0, "llm_events": 0, "parse_fail": 0,
"tokens": 0, "pending_before": 0, "pending_after": 0,
"card_demand": []}
config_error = False
pending = _load_pending()
stats["pending_before"] = len(pending)
fresh: list[dict] = []
stats["card_demand"] = _read_card_demand(args.pipeline_db)
stopped = None
try:
for domain in SRC_DOMAINS:
_walk_domain(ctx, domain, args, ledgers, recents, fresh, stats)
if not args.no_llm and (pending or fresh):
try:
from sanguo_api.llm import resolve_config, LLMClient
cfg = resolve_config()
except Exception as e: # noqa: BLE001 配置缺失
log.error("LLM 配置缺失(env 三件见 runbook P4-1 节): %s", e)
config_error = True
else:
client = LLMClient(cfg)
ctx.reset_stage() # limit=LLM 调用数, 段内预算
try:
n = asyncio.run(_llm_lane(ctx, client, cfg.model,
pending + fresh, ledgers, recents))
stats["llm_events"] = n
except cd.BudgetStop:
stopped = "budget"
ctx.commit()
except cd.WallClockStop:
stopped = "wallclock"
ctx.commit()
except cd.BudgetStop:
stopped = "budget"
ctx.commit()
# 残余回写: 消费 marker 过滤后重排(done 件自然消失, 幂等收敛)
remainder = [r for r in (pending + fresh)
if not cd.is_done("funnel", "llm", r["content_hash"])]
_save_pending(remainder)
stats["pending_after"] = len(remainder)
stats["tokens"] = ctx.tokens_used
stats["llm_calls"] = ctx.stage_units
rc = 1 if config_error else ctx.rc()
print("[funnel] 完成 统计: " + json.dumps(stats, ensure_ascii=False)
+ (f" stopped={stopped}" if stopped else ""))
os.close(fd)
return rc
if __name__ == "__main__":
sys.exit(main())