From a889b2ff0a5f126f9be168af88b81573d6f13a6e Mon Sep 17 00:00:00 2001 From: claude_dev Date: Mon, 7 Sep 2026 18:37:56 +0800 Subject: [PATCH] =?UTF-8?q?fix(data):=20corpus=20=5Ffilter=5Fnew=20?= =?UTF-8?q?=E6=89=B9=E5=86=85=E5=8E=BB=E9=87=8D=E2=80=94=E2=80=94=E5=B7=A8?= =?UTF-8?q?=E6=BD=AE=E5=8D=95=E5=93=8D=E5=BA=94=E5=90=8Cid=E4=B8=A4?= =?UTF-8?q?=E6=9D=A1=E5=85=A8=E6=94=BE=E8=A1=8C=E6=88=90=E9=87=8D=E9=94=AE?= =?UTF-8?q?=E8=A1=8C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 09-07 首跑实锤 ann_meta 3行重键(verify_corpus FAIL): 1225400726/1225487003 整行完全相同×2 + 1225009396 ann_time差1秒(目录挂两条)。根因=has_any只对账 已有键,同一批内重复在统一add前全数放行。修=批内seen截留(首个出现者胜), news_meta复合键同路径受益。事故形态钉进回归测试。 [nas] --- scripts/data_platform/corpus_download.py | 13 ++++++++++-- tests/data_platform/test_corpus_download.py | 23 +++++++++++++++++++++ 2 files changed, 34 insertions(+), 2 deletions(-) diff --git a/scripts/data_platform/corpus_download.py b/scripts/data_platform/corpus_download.py index 47437a4..42a4eb6 100644 --- a/scripts/data_platform/corpus_download.py +++ b/scripts/data_platform/corpus_download.py @@ -672,8 +672,17 @@ def _filter_new(domain, norm_rows, ledger, recent): keys = [_row_key(domain, r) for r in norm_rows] known_l = ledger.has_any(keys) known_r = recent.has_any(keys) - return [r for r, kl, kr in zip(norm_rows, known_l, known_r) - if not kl and not kr] + # 批内去重(09-07 首跑实锤巨潮单响应同 id 两条: 整行相同/时间差1秒): + # has_any 只对账已有键,批内重复须在此截留,否则双双落盘成重键行。 + seen = set() + out = [] + for r, kl, kr in zip(norm_rows, known_l, known_r): + k = _row_key(domain, r) + if kl or kr or k in seen: + continue + seen.add(k) + out.append(r) + return out def _absorb_new(domain, norm_rows, ledgers, recents): diff --git a/tests/data_platform/test_corpus_download.py b/tests/data_platform/test_corpus_download.py index b60bd3e..b3b1386 100644 --- a/tests/data_platform/test_corpus_download.py +++ b/tests/data_platform/test_corpus_download.py @@ -218,6 +218,29 @@ def test_commit_batches_data_before_markers(root, monkeypatch): assert cd.IdLedger("ann_meta").has_any(["m1"]) == [True] +def test_absorb_new_dedups_within_batch(root): + """批内去重(09-07 首跑实锤 ann_meta 3 行重键, verify_corpus FAIL): + 巨潮单次响应同 id 可出现两次——①整行完全相同 ②ann_time 差 1 秒(目录挂两条) + ——都必须只落一行。根因: _filter_new 对整批先算 known 再统一 add,批内重复全放行。""" + import pandas as pd + ledgers = {"ann_meta": cd.IdLedger("ann_meta")} + recents = {"ann_meta": cd.RecentIndex("ann_meta", days=7)} + cd._DAY_BUFFERS.clear() + norm = [cd.norm_ann_row(r, "2026-09-07") for r in ( + _ann_raw(aid="1225400726"), # 完全相同 ×2 + _ann_raw(aid="1225400726"), + _ann_raw(aid="1225009396", ts_ms=1773616809000), # 时间戳差 1 秒 + _ann_raw(aid="1225009396", ts_ms=1773616810000), + _ann_raw(aid="1225487003"), # 正常单条 + )] + new = cd._absorb_new("ann_meta", norm, ledgers, recents) + assert len(new) == 3 + cd.flush_domains() + today = dt.date.today().isoformat() + df = pd.read_parquet(root / "ann_meta" / f"dt={today}" / "part-0.parquet") + assert len(df) == 3 and df.announcement_id.is_unique + + def test_reconcile_ledgers_heals_crash_gap(root): """启动对账自愈(09-07 巡检实锤旧代码丢内存账本 ann 差4.5k/news 差9.9w): 分区有而账本无 → 启动补齐;账本⊇近30日分区永久成立。"""