- system.yaml 新增 llm_scenes(translation / daily_report,含用途/方法/模型要求说明) - load_llm_config(scene=...) 场景覆盖;日报摘要 temperature 0.3 硬编码 → 配置 - M3 去重:唯一篇记录 source_ids(跨源重复合并,首个来源为 source_id) - source_ids 经翻译透传至 events,日报事件 source 多来源拼接展示(≤3 个) - 已部署 pi5:merge 实证 10 源合并;日报 report_id=224 正常入库
264 lines
8.1 KiB
Python
264 lines
8.1 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
|
|
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"))
|
|
articles.append(ProcessedArticle(**data))
|
|
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.ingest(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",
|
|
)
|
|
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)
|
|
except Exception as e:
|
|
logger.exception("来源合并失败 %s: %s", target, e)
|
|
|
|
|
|
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)
|