feat(data): T4 backfill lane ①ann_meta ②news_meta 落地 [nas]
CI/CD / test (push) Successful in 23s
CI/CD / nas-deploy (push) Successful in 26s
CI/CD / nas-verify (push) Successful in 6s

spec §18.7 T4(unit marker 断点+--limit 冒烟不落marker+墙钟rc3 全链复用):
- ①ann_meta 回补: unit=(stock,year) 2000→今年逐年按池序, 空年=合法done
  (探测优先);今年窗口右端=today 相对;~15万 unit @1s≈3-6 跑日
- ②news_meta 回补: np-listapi unit=(stock,page) 翻到空页止;空页+page≥200
  双止条件都落 news_end 哨兵(到头股零重扫);mkt前缀 6→1/其余→0
- 共用 _absorb_new(daily/backfill 同一条 id 去重+追加管道)
- 测试 46 个;修 wiring 网络泄漏:_wire_backfill_ann 曾漏锁
  fetch_news_backfill_page 致 ann 测试真打 np-listapi(40s/真出网,
  零网络铁律违反)——补锁后套件 159s→1.3s;全套 263 绿
This commit is contained in:
2026-09-06 23:19:50 +08:00
parent e22f0e0aa8
commit 10c486c64d
2 changed files with 257 additions and 13 deletions
+106 -13
View File
@@ -536,6 +536,17 @@ def _filter_new(domain, norm_rows, ledger, recent):
if not kl and not kr]
def _absorb_new(domain, norm_rows, ledgers, recents):
"""过滤已知 id → 追加分区+账本+近窗(daily/backfill 共用)。返回新行。"""
new = _filter_new(domain, norm_rows, ledgers[domain], recents[domain])
if new:
append_parquet(domain, new, ID_KEYS[domain])
keys = [_row_key(domain, r) for r in new]
ledgers[domain].add(keys)
recents[domain].add(keys)
return new
def _record_depth(ctx, source, code, date_str):
if not date_str:
return
@@ -598,12 +609,8 @@ def _run_ann_increment(ctx, pool, cl_cn, ledgers, recents):
continue
candidates.extend(collect_pdf_candidates(rows))
norm = [norm_ann_row(r, first_seen) for r in rows]
new = _filter_new("ann_meta", norm, ledgers["ann_meta"], recents["ann_meta"])
new = _absorb_new("ann_meta", norm, ledgers, recents)
if new:
append_parquet("ann_meta", new, "announcement_id")
keys = [_row_key("ann_meta", r) for r in new]
ledgers["ann_meta"].add(keys)
recents["ann_meta"].add(keys)
for r in new:
_record_depth(ctx, "ann", r["sec_code"] or code,
(r["ann_time"] or "")[:10])
@@ -635,12 +642,8 @@ def _run_news_increment(ctx, pool, cl_em, ledgers, recents, new_articles):
log.warning("news unit %s 失败(不标done): %s", code, e)
continue
norm = [norm_news_search_row(r, code) for r in rows]
new = _filter_new("news_meta", norm, ledgers["news_meta"], recents["news_meta"])
new = _absorb_new("news_meta", norm, ledgers, recents)
if new:
append_parquet("news_meta", new, ["art_code", "stock_code"])
keys = [_row_key("news_meta", r) for r in new]
ledgers["news_meta"].add(keys)
recents["news_meta"].add(keys)
for r in new:
_record_depth(ctx, "news", code, (r["show_time"] or "")[:10])
new_articles.append((r["art_code"], r["url"], r["show_time"]))
@@ -778,11 +781,101 @@ def run_daily(ctx, pool):
led.flush()
# ---------- backfill lane(T2 占位, ①②=T4 ③④=T5) ----------
# ---------- backfill lane(①ann_meta ②news_meta;③④=T5) ----------
def _run_ann_backfill(ctx, pool, cl_cn, ledgers, recents):
"""①公告元数据回补 2000→今, unit=(stock,year), 逐年按池序(探测优先:空年=合法)。"""
today = dt.date.today()
first_seen = today.isoformat()
since_flush = 0
for year in range(BACKFILL_FROM_YEAR, today.year + 1):
year_end = today.isoformat() if year == today.year else f"{year}-12-31"
for code, org in pool:
unit = f"{code}_{year}"
if ctx.limit is not None and ctx.units >= ctx.limit:
return
if is_done(ctx.lane, "ann", unit):
continue
try:
rows = fetch_cninfo_announcements(cl_cn, code, org,
f"{year}-01-01", year_end)
except DomainCooldown as e:
log.warning("backfill ann 域冷却 @%s: %s", unit, e)
ctx.rate_limited |= e.rate_limited
ctx.hard_cool |= not e.rate_limited
return
except (TransportError, HttpDeterministicError) as e:
ctx.failed += 1
log.warning("backfill ann %s 失败(不标done): %s", unit, e)
continue
norm = [norm_ann_row(r, first_seen) for r in rows]
new = _absorb_new("ann_meta", norm, ledgers, recents)
for r in new:
_record_depth(ctx, "ann", r["sec_code"] or code,
(r["ann_time"] or "")[:10])
if new:
log.info("backfill ann %s: +%d/%d", unit, len(new), len(rows))
if ctx.limit is None:
mark_done(ctx.lane, "ann", unit)
ctx.units += 1
since_flush = _maybe_flush(ledgers, since_flush + 1)
ctx.stop_now()
def _run_news_backfill(ctx, pool, cl_em, ledgers, recents):
"""②新闻标题级回补(np-listapi), unit=(stock,page), 翻到空页或 page≥200 止。"""
for code, _org in pool:
if is_done(ctx.lane, "news_end", code):
continue
mkt = "1" if code.startswith("6") else "0"
capped = False
for page in range(1, NEWS_BACKFILL_MAX_PAGE + 1):
unit = f"{code}_p{page}"
if ctx.limit is not None and ctx.units >= ctx.limit:
return
if is_done(ctx.lane, "news", unit):
continue
try:
rows = fetch_news_backfill_page(cl_em, code, mkt, page)
except DomainCooldown as e:
log.warning("backfill news 域冷却 @%s: %s", unit, e)
ctx.rate_limited |= e.rate_limited
ctx.hard_cool |= not e.rate_limited
return
except (TransportError, HttpDeterministicError) as e:
ctx.failed += 1
log.warning("backfill news %s 失败(断点留此页): %s", unit, e)
break
if not rows:
if ctx.limit is None:
mark_done(ctx.lane, "news", unit)
mark_done(ctx.lane, "news_end", code)
log.info("backfill news %s: 到头(空页@%d)", code, page)
break
norm = [norm_news_np_row(r, code) for r in rows]
new = _absorb_new("news_meta", norm, ledgers, recents)
for r in new:
_record_depth(ctx, "news", code, (r["show_time"] or "")[:10])
if ctx.limit is None:
mark_done(ctx.lane, "news", unit)
ctx.units += 1
if page == NEWS_BACKFILL_MAX_PAGE:
capped = True
ctx.stop_now()
if capped and ctx.limit is None:
mark_done(ctx.lane, "news_end", code)
log.info("backfill news %s: page=%d 硬顶止", code, NEWS_BACKFILL_MAX_PAGE)
def run_backfill(ctx, pool):
log.warning("backfill lane 未实施(T4/T5): 让路退出, daily 数据不受影响")
return
ledgers = {d: IdLedger(d) for d in ID_KEYS}
recents = {d: RecentIndex(d) for d in ID_KEYS}
cl_cn = DomainClient("cninfo")
cl_em = DomainClient("eastmoney")
_run_ann_backfill(ctx, pool, cl_cn, ledgers, recents)
_run_news_backfill(ctx, pool, cl_em, ledgers, recents)
for led in ledgers.values():
led.flush()
# ---------- lane 入口 ----------