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日分区永久成立。"""