diff --git a/scripts/data_platform/corpus_download.py b/scripts/data_platform/corpus_download.py index 6d7bb6c..05b6e26 100644 --- a/scripts/data_platform/corpus_download.py +++ b/scripts/data_platform/corpus_download.py @@ -872,6 +872,25 @@ def _run_pdf_daily(ctx, pool, cl_cn, ledgers, candidates): _append_pdf_index(index_rows) +def _reconcile_ledgers(ledgers, recents): + """启动对账: 近窗分区有而账本无的 key 补进账本(一次性自愈)。 + + 09-07 巡检实锤: 旧代码部分段无 flush 计数器, docker stop 的腿丢在内存 + 账本(ann 差 4.5k/news 差 9.9w, 分区完好);此后任何崩溃窗(批量提交周期 + 内)同理。对账后账本 ⊇ 近30日分区, 永久闭合「分区有/账本无」缺口。 + """ + for domain, recent in recents.items(): + keys = list(recent.keys) + if not keys: + continue + known = ledgers[domain].has_any(keys) + missing = [k for k, has in zip(keys, known) if not has] + if missing: + ledgers[domain].add(missing) + ledgers[domain].flush() + log.info("ledger reconcile %s: +%d (近窗分区对账)", domain, len(missing)) + + def run_daily(ctx, pool): ledgers = {d: IdLedger(d) for d in ID_KEYS} recents = {d: RecentIndex(d) for d in ID_KEYS} @@ -880,6 +899,7 @@ def run_daily(ctx, pool): cl_art = DomainClient("article") ctx.clients = {"cninfo": cl_cn, "eastmoney": cl_em, "article": cl_art} ctx.set_stores(ledgers) + _reconcile_ledgers(ledgers, recents) candidates = _run_ann_increment(ctx, pool, cl_cn, ledgers, recents) _run_news_increment(ctx, pool, cl_em, ledgers, recents, new_articles := []) _run_fulltext(ctx, cl_art, ledgers, recents, new_articles) @@ -1111,6 +1131,7 @@ def run_backfill(ctx, pool): cl_art = DomainClient("article") ctx.clients = {"cninfo": cl_cn, "eastmoney": cl_em, "article": cl_art} ctx.set_stores(ledgers) + _reconcile_ledgers(ledgers, recents) _run_ann_backfill(ctx, pool, cl_cn, ledgers, recents) _run_news_backfill(ctx, pool, cl_em, ledgers, recents) _run_fulltext_backfill(ctx, cl_art, ledgers, recents) diff --git a/tests/data_platform/test_corpus_download.py b/tests/data_platform/test_corpus_download.py index 3ffc71a..b60bd3e 100644 --- a/tests/data_platform/test_corpus_download.py +++ b/tests/data_platform/test_corpus_download.py @@ -218,6 +218,24 @@ def test_commit_batches_data_before_markers(root, monkeypatch): assert cd.IdLedger("ann_meta").has_any(["m1"]) == [True] +def test_reconcile_ledgers_heals_crash_gap(root): + """启动对账自愈(09-07 巡检实锤旧代码丢内存账本 ann 差4.5k/news 差9.9w): + 分区有而账本无 → 启动补齐;账本⊇近30日分区永久成立。""" + import pandas as pd + today = dt.date.today().isoformat() + d = root / "ann_meta" / f"dt={today}" + d.mkdir(parents=True) + pd.DataFrame({"announcement_id": ["r1", "r2"]}).to_parquet(d / "part-0.parquet") + ledgers = {"ann_meta": cd.IdLedger("ann_meta")} # 空账本 + recents = {"ann_meta": cd.RecentIndex("ann_meta", days=7)} + assert ledgers["ann_meta"].has_any(["r1", "r2"]) == [False, False] + cd._reconcile_ledgers(ledgers, recents) + assert cd.IdLedger("ann_meta").has_any(["r1", "r2"]) == [True, True] # 落盘可复读 + # 幂等: 再对账零新增 + cd._reconcile_ledgers(ledgers, recents) + assert cd.IdLedger("ann_meta").has_any(["r1", "r2"]) == [True, True] + + # ---------- unit marker ---------- def test_marker_done_and_skip(root):