From 95b1be2c2ef26b427c39829e5bc85bc31574aa21 Mon Sep 17 00:00:00 2001 From: butubb <1422726308@qq.com> Date: Sat, 10 Oct 2026 17:24:08 +0800 Subject: [PATCH] =?UTF-8?q?fix(monitor):=20=E4=B8=80=E8=BD=AE=E4=BA=A7?= =?UTF-8?q?=E7=89=A9=E9=87=8C=E9=87=8D=E5=A4=8D=E7=9A=84=E4=BD=9C=E5=93=81?= =?UTF-8?q?=E4=BC=9A=E8=AE=A9=E6=8C=87=E6=A0=87=E5=BF=AB=E7=85=A7=E6=92=9E?= =?UTF-8?q?=E5=94=AF=E4=B8=80=E9=94=AE=EF=BC=8C=E6=95=B4=E4=B8=AA=20run=20?= =?UTF-8?q?=E5=B4=A9=E6=8E=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 跑真任务时踩到的: sqlalchemy.exc.IntegrityError: (1062, "Duplicate entry '7-7690458980574358513-45' for key 'uq_note_metric'") 根因是我上一版写错了一处作用域:退化路径里那个「遍历已知作品」的循环写在了**目标循环内部**, 所以任务有多个目标时,同一批已知作品会被拉两遍 → 同一件作品在一轮里出现两条记录 → ingest 给同一件作品写两份本轮快照 → 撞 (task_id, note_id, run_id) 唯一键。 两处都修,各挡一层: * douyin_fetch:去重集合挪到 collect 的最外层,**跨目标**只算一次;退化时也先查一遍 已知作品是否已刷过。 * ingest:`_ingest_notes` 对「一轮里重复出现的 note_id」免疫。一层在源头、一层在入口, 因为产物里重复并不罕见(多个目标指向同一个人、上游重跑、退化路径),不该靠上游自觉。 测试 +2:多目标时已知作品只刷一次;同一轮里重复的作品只落一份快照(这条会崩在 唯一键上,所以它测的正是运行时的那个崩法)。 --- api/monitor/douyin_fetch.py | 8 +++++++- api/monitor/ingest.py | 7 +++++++ tests/test_douyin_fetch.py | 31 +++++++++++++++++++++++++++++++ tests/test_monitor_ingest.py | 22 ++++++++++++++++++++++ 4 files changed, 67 insertions(+), 1 deletion(-) diff --git a/api/monitor/douyin_fetch.py b/api/monitor/douyin_fetch.py index 9dbe3a0..0f4949c 100644 --- a/api/monitor/douyin_fetch.py +++ b/api/monitor/douyin_fetch.py @@ -67,9 +67,13 @@ async def collect( comments: List[Dict[str, Any]] = [] errors: List[str] = [] + # **整个 collect 只去重一次的、跨目标的集合**:退化路径会把「库里已知的全部作品」 + # 在每个目标下都刷一遍,多个目标就会出现同一件作品好几条记录 —— 而一对一快照的 + # 唯一键是 (task_id, note_id, run_id),同一条作品在一轮里出现两次会直接撞键。 + seen_aweme: set = set() + for target in targets: sec_user_id = target.external_id - seen_aweme: set = set() # 作品列表是主路径;它被那道真校验挡着时,退化到「标题/昵称靠主页接口,作品靠 # 已知 id 逐条刷新」—— 拿不到新作品,但已知作品的指标还能继续更新。 @@ -81,6 +85,8 @@ async def collect( errors.append(f"拉取博主 {sec_user_id} 的作品列表失败:{exc}") videos = [] for aweme_id in known_aweme_ids: + if aweme_id in seen_aweme: + continue try: videos.append(await douyin_api.video_detail(aweme_id, cookie=cookie)) except douyin_api.DouyinApiError as detail_exc: diff --git a/api/monitor/ingest.py b/api/monitor/ingest.py index 79428f4..b6a12f2 100644 --- a/api/monitor/ingest.py +++ b/api/monitor/ingest.py @@ -341,11 +341,18 @@ async def _ingest_notes( """ now = get_current_timestamp() new_count = 0 + # 同一轮里重复出现的作品只处理一次。**指标快照的唯一键是 (task_id, note_id, run_id)**, + # 同一件作品在一轮里进来两次会让第二次插入直接撞键、整个 run 崩掉 —— 产物里重复并不 + # 罕见(多个目标指向同一个人、或退化路径重复刷新)。 + seen_in_run: set = set() for record in records: note_id = adapter.note_field(record, "note_id") if not note_id: continue + if note_id in seen_in_run: + continue + seen_in_run.add(note_id) note = await session.scalar( select(MonitorNote).where( diff --git a/tests/test_douyin_fetch.py b/tests/test_douyin_fetch.py index 12b324a..fa5dbe3 100644 --- a/tests/test_douyin_fetch.py +++ b/tests/test_douyin_fetch.py @@ -158,6 +158,37 @@ class TestDegradation: # 但错误照样报出来 —— 这一轮是「部分可用」,不是「一切正常」,别粉饰。 assert any("作品列表失败" in error for error in result["errors"]) + @pytest.mark.asyncio + async def test_known_works_are_not_refetched_for_every_target(self, monkeypatch, tmp_path): + """退化路径不能每个目标都把同一批已知作品再刷一遍。 + + 任务有多个目标时那会让同一件作品在一轮里出现两次,而指标快照的唯一键是 + (task_id, note_id, run_id) —— 第二次插入直接撞键,整个 run 崩掉(踩过: + Duplicate entry for key 'uq_note_metric')。 + """ + + async def blocked(sec_user_id, count=20, *, cookie=""): + raise douyin_api.DouyinApiError("接口返回了空内容") + + calls = [] + + async def fake_detail(aweme_id, *, cookie=""): + calls.append(aweme_id) + return _video(aweme_id) + + monkeypatch.setattr(douyin_api, "author_videos", blocked) + monkeypatch.setattr(douyin_api, "video_detail", fake_detail) + + result = await _collect( + tmp_path, + want_comments=False, + targets=[_Target("sec-a"), _Target("sec-b")], + known_aweme_ids=["999"], + ) + + assert calls == ["999"], "同一件已知作品只该刷一次,而不是每个目标一次" + assert result["notes"] == 1 + @pytest.mark.asyncio async def test_nothing_at_all_still_reports_the_reason(self, monkeypatch, tmp_path): async def blocked(sec_user_id, count=20, *, cookie=""): diff --git a/tests/test_monitor_ingest.py b/tests/test_monitor_ingest.py index 7bccd78..84bd864 100644 --- a/tests/test_monitor_ingest.py +++ b/tests/test_monitor_ingest.py @@ -725,6 +725,28 @@ class TestDouyinIngest: assert len(await _events(db, EVENT_METRIC_DELTA)) == 1 +class TestDuplicateRecordsInOneRun: + """一轮产物里重复出现的作品只该算一次。 + + 指标快照的唯一键是 ``(task_id, note_id, run_id)``:同一条作品在一轮里进来两次, + 第二次插入会撞键并让整个 run 崩掉 —— 产物里重复并不罕见(多个目标指向同一个人、 + 或退化路径重复刷新)。 + """ + + @pytest.mark.asyncio + async def test_a_duplicated_note_is_processed_once(self, db, tmp_path): + task = await _make_task(db) + _write_run_dir(tmp_path, [_note("n1"), _note("n1")]) + run = await _make_run(db, task, started_at=1) + + result = await ingest_run(db, run, task, tmp_path) + + assert result.new_notes == 1 + assert len(list((await db.scalars(select(MonitorNote))).all())) == 1 + # 崩就崩在这一句上:一条作品只能有一份本轮快照。 + assert len(list((await db.scalars(select(MonitorNoteMetric))).all())) == 1 + + class TestNicknameRefresh: """已入库的评论,昵称要跟着重新采集的值走。