diff --git a/app/cli.py b/app/cli.py index 53361be..c7b3bcf 100644 --- a/app/cli.py +++ b/app/cli.py @@ -312,6 +312,26 @@ def report(): raise typer.Exit(code=1) +@app.command() +def sync_sources(): + """M3/M9 辅助:把 deduped 已合并的多来源同步到 events/Qdrant/MySQL""" + from dedup.pipeline import sync_sources_from_deduped + + try: + stats = sync_sources_from_deduped() + typer.echo( + f"\n✅ 来源同步完成: 扫描 {stats['deduped_scanned']} 个唯一篇, " + f"多来源 {stats['multi_source']} 个, " + f"events 更新 {stats['events_updated']}, " + f"Qdrant 更新 {stats['qdrant_updated']}, " + f"MySQL 更新 {stats['db_updated']}" + ) + except Exception as e: + logger.exception("来源同步失败") + typer.echo(f"❌ 来源同步出错: {e}", err=True) + raise typer.Exit(code=1) + + @app.command() def pipeline( skip_report: bool = typer.Option( diff --git a/dedup/pipeline.py b/dedup/pipeline.py index 4892b36..628ef5a 100644 --- a/dedup/pipeline.py +++ b/dedup/pipeline.py @@ -172,10 +172,160 @@ def _merge_duplicate_source(article: ProcessedArticle, result: DedupResult) -> N ) logger.debug("来源合并: %s → %s (sources=%s)", article.source_id, result.matched_url_hash, merged) + # 同步到已生成的 events / Qdrant,避免“deduped 有多来源但 DB 没有” + _sync_merged_sources_downstream(result.matched_url_hash, merged) except Exception as e: logger.exception("来源合并失败 %s: %s", target, e) +def _sync_event_source_ids(url_hash: str, source_ids: list[str]) -> int: + """把合并后的 source_ids 同步到已生成的 events JSON。""" + updated = 0 + for ev_path in sorted(Path("data/events").glob(f"*/{url_hash}.json")): + try: + data = json.loads(ev_path.read_text(encoding="utf-8")) + if data.get("source_ids") != source_ids: + data["source_ids"] = list(source_ids) + ev_path.write_text( + json.dumps(data, indent=2, ensure_ascii=False), + encoding="utf-8", + ) + logger.info("来源合并同步 events: %s -> %s", ev_path, source_ids) + updated += 1 + except Exception as e: + logger.exception("同步 events 来源失败 %s: %s", ev_path, e) + return updated + + +def _sync_qdrant_source_ids(url_hash: str, source_ids: list[str]) -> int: + """把合并后的 source_ids 同步到 Qdrant 中已存在的向量点。""" + emb_paths = sorted(Path("data/embeddings").glob(f"*/{url_hash}.json")) + if not emb_paths: + return 0 + + event_path = next(Path("data/events").glob(f"*/{url_hash}.json"), None) + if event_path is None: + return 0 + + try: + article_data = json.loads(event_path.read_text(encoding="utf-8")) + except Exception as e: + logger.warning("读取 events 失败,跳过 Qdrant 来源同步 %s: %s", event_path, e) + return 0 + + try: + from vectorstore.client import VectorStore, make_qdrant_client + from vectorstore.pipeline import _build_payload + + client = make_qdrant_client() + store = VectorStore(client) + try: + payload = _build_payload(article_data) + for emb_path in emb_paths: + emb_data = json.loads(emb_path.read_text(encoding="utf-8")) + store.upsert([{ + "id": url_hash, + "vector": emb_data["vector"], + "payload": payload, + }]) + finally: + store.close() + logger.info("来源合并同步 Qdrant: %s -> %s", url_hash, source_ids) + return len(emb_paths) + except Exception as e: + logger.warning("同步 Qdrant 来源失败 %s: %s", url_hash, e) + return 0 + + +def _sync_merged_sources_downstream(url_hash: str, source_ids: list[str]) -> None: + """去重合并来源后,同步更新 events 与 Qdrant 中已存在的下游数据。""" + _sync_event_source_ids(url_hash, source_ids) + _sync_qdrant_source_ids(url_hash, source_ids) + + +def sync_sources_from_deduped() -> dict: + """历史数据回填:扫描所有 deduped 唯一篇,把多来源同步到 events/Qdrant/MySQL。 + + 该函数用于修复“deduped 已合并多来源,但旧 events/数据库未更新”的历史数据。 + """ + stats = { + "deduped_scanned": 0, + "multi_source": 0, + "events_updated": 0, + "qdrant_updated": 0, + "db_updated": 0, + } + + targets = sorted(Path("data/deduped").glob("*/uniques/*.json")) + for target in targets: + try: + data = json.loads(target.read_text(encoding="utf-8")) + except Exception: + continue + stats["deduped_scanned"] += 1 + + raw_ids = data.get("source_ids") or [data.get("source_id")] + source_ids = list(dict.fromkeys(x for x in raw_ids if x)) + if len(source_ids) <= 1: + continue + + stats["multi_source"] += 1 + url_hash = data.get("url_hash", target.stem) + stats["events_updated"] += _sync_event_source_ids(url_hash, source_ids) + stats["qdrant_updated"] += _sync_qdrant_source_ids(url_hash, source_ids) + stats["db_updated"] += _update_db_sources(url_hash, source_ids) + + logger.info( + "来源回填完成: scanned=%d multi=%d events=%d qdrant=%d db=%d", + stats["deduped_scanned"], stats["multi_source"], + stats["events_updated"], stats["qdrant_updated"], stats["db_updated"], + ) + return stats + + +def _update_db_sources(url_hash: str, source_ids: list[str]) -> int: + """根据最新 events 来源,更新 MySQL news_event 中该文章的 source/sources。""" + event_path = next(Path("data/events").glob(f"*/{url_hash}.json"), None) + if event_path is None: + return 0 + try: + article_data = json.loads(event_path.read_text(encoding="utf-8")) + except Exception: + return 0 + + try: + from scheduler.reporter import _article_source_label, _article_sources_full + source_label = _article_source_label(article_data) + source_list = _article_sources_full(article_data) + if not source_label and not source_list: + return 0 + + from report_db import connect + conn = connect() + try: + with conn.cursor() as cur: + cur.execute( + """ + UPDATE news_event + SET source = %s, sources = %s + WHERE section = 'intl' AND url = %s + """, + ( + source_label, + json.dumps(source_list, ensure_ascii=False) if source_list else None, + article_data.get("url"), + ), + ) + updated = cur.rowcount + conn.commit() + return updated or 0 + finally: + conn.close() + except Exception as e: + logger.warning("同步 MySQL 来源失败 %s: %s", url_hash, e) + return 0 + + def _layer_num(result: DedupResult) -> int: """DedupResult → 命中层编号。""" if result.matched_layer is None: diff --git a/docs/pipeline.md b/docs/pipeline.md index b79fd8f..a0ed7f9 100644 --- a/docs/pipeline.md +++ b/docs/pipeline.md @@ -45,6 +45,10 @@ uv run en-news extract --source cnbc 2. 正文 hash 一致 → 重复 3. SimHash 汉明距离 ≤ 阈值 → 模糊重复 - 跨源合并:重复篇的来源 ID 会追加到保留唯一篇的 `source_ids`,日报可展示多来源。 +- 合并后会同步更新已生成的 `events` 与 Qdrant payload;历史数据可执行: + ```bash + uv run en-news sync-sources + ``` ```bash uv run en-news dedup diff --git a/docs/usage.md b/docs/usage.md index 12b6b68..0ada3d8 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -12,6 +12,7 @@ | `index` | M6 写入 Qdrant | `--recreate` 重建;`--date/-d` 指定日期;`--all` 全量回灌 | | `search` | M6 语义检索 | 必需 query;`--top-k` 返回条数 | | `report` | M9 生成日报并入库 | 无 | +| `sync-sources` | 把 deduped 已合并的多来源同步到 events/Qdrant/MySQL | 无 | | `pipeline` | M2→M6(+日报) | `--skip-report` 跳过日报;`--date/-d` 指定日期 | | `mcp-server` | M8 启动 MCP 服务 | 无 | @@ -38,6 +39,9 @@ uv run en-news index --date 20260801 # 全量回灌历史向量(--recreate 搭配 --all 时只会在首个日期重建 collection) uv run en-news index --all uv run en-news embed --all + +# 修复历史多来源:把 deduped 中已合并的 source_ids 同步到 events/Qdrant/MySQL +uv run en-news sync-sources ``` ## 2. Shell 脚本