Files
intl_news/dedup/pipeline.py
T
simon 9e4f5b4c75 fix: 去重合并多来源后同步到 events/Qdrant/MySQL,新增 sync-sources 回填命令
- dedup/pipeline: 合并 source_ids 后自动同步已生成的 events JSON 和 Qdrant payload
- 新增 sync_sources_from_deduped / uv run en-news sync-sources,用于历史数据回填
- sync-sources 会扫描 deduped 多来源唯一篇,更新 events、Qdrant、MySQL news_event.sources
- docs: 补充 sync-sources 使用说明
2026-08-23 10:55:23 +08:00

423 lines
14 KiB
Python

"""批量去重管道:扫描 processed 目录 → 判重 → 唯一条目写入 deduped。
输入: data/processed/{source_id}/{YYYYMMDD}/{url_hash}.json
输出: data/deduped/{YYYYMMDD}/uniques/{url_hash}.json
"""
import json
import logging
from datetime import datetime
from pathlib import Path
from crawler.utils import get_news_day
from dedup.deduper import Deduper, article_to_fingerprint
from dedup.models import DedupResult
from extractor.models import ProcessedArticle
logger = logging.getLogger(__name__)
def _get_processed_sources(base_dir: str = "data/processed") -> list[str]:
"""扫描 data/processed/ 下所有源 ID。
Args:
base_dir: processed 数据根目录
Returns:
源 ID 列表
"""
raw_path = Path(base_dir)
if not raw_path.exists():
return []
return sorted([
d.name for d in raw_path.iterdir()
if d.is_dir() and not d.name.startswith(".")
])
def _load_processed_articles(
source_id: str,
date_str: str,
) -> list[ProcessedArticle]:
"""加载指定源/日期的已处理文章。
Args:
source_id: 新闻源 ID
date_str: 日期 YYYYMMDD
Returns:
ProcessedArticle 列表
"""
base_dir = Path(f"data/processed/{source_id}/{date_str}")
if not base_dir.exists():
return []
articles: list[ProcessedArticle] = []
for json_file in sorted(base_dir.glob("*.json")):
# 跳过 index.jsonl
if json_file.name == "index.jsonl":
continue
try:
data = json.loads(json_file.read_text(encoding="utf-8"))
article = ProcessedArticle(**data)
# M2 可能写入 no_content / failed:这些文章没有有效正文,
# 不应进入去重、翻译、向量化等下游环节。
if article.status != "success" or not article.content.strip():
continue
articles.append(article)
except (json.JSONDecodeError, Exception) as e:
logger.warning("解析 processed JSON 失败 %s: %s", json_file, e)
return articles
def dedup_source(
source_id: str,
deduper: Deduper,
date_str: str | None = None,
) -> dict:
"""对单个源的已处理文章执行去重。
Args:
source_id: 新闻源 ID
deduper: 去重器实例
date_str: 日期 YYYYMMDD,默认当前新闻日
Returns:
统计 dict
"""
if date_str is None:
date_str = get_news_day()
logger.info("━━━ 去重 [%s] %s ━━━", source_id, date_str)
articles = _load_processed_articles(source_id, date_str)
if not articles:
logger.warning("[%s] %s 无待处理文章", source_id, date_str)
return {"source_id": source_id, "total": 0, "unique": 0, "duplicate": 0}
# 输出目录
out_dir = Path(f"data/deduped/{date_str}/uniques")
out_dir.mkdir(parents=True, exist_ok=True)
unique_count = 0
dup_count = 0
for article in articles:
# 先只读判断,不写指纹;等唯一文件落盘成功后再写指纹,
# 避免“指纹已入库但唯一文件未生成”导致该文章后续被当重复丢弃。
result = deduper.check(article)
if result.is_duplicate:
dup_count += 1
# 跨源重复:把来源合并进已保留的唯一篇(记录多个来源)
_merge_duplicate_source(article, result)
logger.debug("[%s] 🔁 %s → L%d: %s",
source_id,
article.title[:40],
_layer_num(result),
result.short_summary())
else:
unique_count += 1
# 初始化来源列表(首个来源 = 本篇文章来源)
if not article.source_ids:
article.source_ids = [article.source_id]
# 写入唯一条目
out_file = out_dir / f"{article.url_hash}.json"
out_file.write_text(
article.model_dump_json(indent=2, ensure_ascii=False),
encoding="utf-8",
)
# 唯一文件落盘成功后再写指纹
deduper.store.upsert(article_to_fingerprint(article))
logger.debug("[%s] ✅ %s (%d words)",
source_id, article.title[:40], article.word_count)
logger.info("[%s] 去重完成: 唯一 %d / 重复 %d / 总计 %d",
source_id, unique_count, dup_count, len(articles))
return {
"source_id": source_id,
"total": len(articles),
"unique": unique_count,
"duplicate": dup_count,
}
def _merge_duplicate_source(article: ProcessedArticle, result: DedupResult) -> None:
"""重复篇:把来源 ID 追加进已保留的唯一篇 JSON(最终显示的新闻记录多个来源)。
唯一篇文件按 url_hash 定位(跨日期目录搜索,因指纹窗口为 ±30 天);
文件不存在(超窗口被清理)时仅记录日志,不阻塞去重流程。
"""
if not result.matched_url_hash:
return
candidates = sorted(Path("data/deduped").glob(f"*/uniques/{result.matched_url_hash}.json"))
if not candidates:
logger.warning(
"重复篇唯一文件不存在(可能已超窗口): %s(重复来源 %s 未合并)",
result.matched_url_hash, article.source_id,
)
return
target = candidates[0]
try:
data = json.loads(target.read_text(encoding="utf-8"))
# 旧格式文件可能无 source_ids:以主来源 source_id 兜底
merged = list(dict.fromkeys(
[*(data.get("source_ids") or [data.get("source_id")]), article.source_id]
))
data["source_ids"] = merged
target.write_text(
json.dumps(data, indent=2, ensure_ascii=False), encoding="utf-8"
)
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:
return 0
mapping = {"url": 1, "content": 2, "simhash": 3}
return mapping.get(result.matched_layer.value, 0)
def dedup_all_sources(
source_filter: str | None = None,
date_str: str | None = None,
) -> dict:
"""对所有源的已处理文章执行去重。
Args:
source_filter: 可选,只处理指定源
date_str: 日期,默认当前新闻日
Returns:
统计 dict
"""
if date_str is None:
date_str = get_news_day()
start_time = datetime.now()
if source_filter:
sources = [source_filter] if source_filter in _get_processed_sources() else []
else:
sources = _get_processed_sources()
logger.info("══════ 开始去重 %d 个源,日期: %s ══════", len(sources), date_str)
total_unique = 0
total_dup = 0
total_articles = 0
with Deduper() as deduper:
for src in sources:
result = dedup_source(src, deduper, date_str)
total_articles += result["total"]
total_unique += result["unique"]
total_dup += result["duplicate"]
# 输出 dedup 索引
_write_dedup_index(deduper, date_str, total_unique)
elapsed = (datetime.now() - start_time).total_seconds()
logger.info("══════ 去重完成: 唯一 %d / 重复 %d / 总计 %d,耗时 %.1f 秒 ══════",
total_unique, total_dup, total_articles, elapsed)
return {
"sources_processed": len(sources),
"total_articles": total_articles,
"unique": total_unique,
"duplicate": total_dup,
"elapsed_sec": elapsed,
"date": date_str,
}
def _write_dedup_index(
deduper: Deduper,
date_str: str,
unique_count: int,
) -> None:
"""写出去重索引文件。
Args:
deduper: 去重器实例
date_str: 日期
unique_count: 唯一文章数
"""
out_dir = Path(f"data/deduped/{date_str}")
out_dir.mkdir(parents=True, exist_ok=True)
stats = deduper.stats()
index_data = {
"date": date_str,
"unique_articles": unique_count,
"fingerprint_db_total": stats.total,
"fingerprint_db_by_source": stats.by_source,
"fingerprint_db_earliest": stats.earliest,
"fingerprint_db_latest": stats.latest,
"generated_at": datetime.now().isoformat(),
}
index_path = out_dir / "index.json"
index_path.write_text(
json.dumps(index_data, indent=2, ensure_ascii=False),
encoding="utf-8",
)
logger.info("去重索引已写入: %s", index_path)