fetch 改为先抓取后回退缓存:成功采集就不用旧缓存,抓取失败才回退 collected_keywords.json
This commit is contained in:
+30
-20
@@ -2,6 +2,10 @@
|
|||||||
|
|
||||||
按 config.sources 启用各可插拔数据源,汇总统一格式行。
|
按 config.sources 启用各可插拔数据源,汇总统一格式行。
|
||||||
单源失败不影响其它源(内部逐个 try),整体再套 with_fallback 兜底。
|
单源失败不影响其它源(内部逐个 try),整体再套 with_fallback 兜底。
|
||||||
|
|
||||||
|
缓存策略(用户要求:成功采集就不用缓存):
|
||||||
|
1. 先用种子词走数据源抓取(Google Trends 等),成功即用新数据;
|
||||||
|
2. 抓取失败/无结果才回退 output/<国>/collected_keywords.json 旧缓存,保证流水线不中断。
|
||||||
"""
|
"""
|
||||||
from typing import Any, Dict, List
|
from typing import Any, Dict, List
|
||||||
|
|
||||||
@@ -18,26 +22,7 @@ def fetch_node(state: Dict[str, Any]) -> Dict[str, Any]:
|
|||||||
rows: List[Dict[str, Any]] = []
|
rows: List[Dict[str, Any]] = []
|
||||||
errors = list(state.get("errors") or [])
|
errors = list(state.get("errors") or [])
|
||||||
|
|
||||||
# 采集缓存优先:采集(fetch_keywords)成功后写入 output/<国>/collected_keywords.json,
|
# 1) 先尝试数据源抓取(用种子词),成功即用新数据
|
||||||
# 这里直接用(跳过 Google 重抓),避免重复撞限流;无缓存才走数据源抓取
|
|
||||||
use_collected = (config.get("fetch") or {}).get("use_collected", True)
|
|
||||||
if use_collected:
|
|
||||||
try:
|
|
||||||
import json as _json
|
|
||||||
from pathlib import Path as _Path
|
|
||||||
p = _Path(state.get("cache_dir") or state.get("output_dir", "")) / "collected_keywords.json"
|
|
||||||
if p.exists():
|
|
||||||
data = _json.loads(p.read_text(encoding="utf-8"))
|
|
||||||
cached_rows = data.get("keywords") or []
|
|
||||||
if cached_rows:
|
|
||||||
rows = [dict(r) for r in cached_rows] # 已过滤去重的关键词
|
|
||||||
print(f"[fetch] 使用采集缓存 {len(rows)} 条({country},跳过 Google 抓取)")
|
|
||||||
stats = dict(state.get("stats") or {})
|
|
||||||
stats["fetch"] = {"raw_rows": len(rows), "sources": ["collected_cache"], "errors": 0}
|
|
||||||
return {"raw_rows": rows, "errors": errors, "stats": stats}
|
|
||||||
except Exception as e: # noqa: BLE001
|
|
||||||
print(f"[fetch] 读取采集缓存失败(回退数据源): {e}")
|
|
||||||
|
|
||||||
for name in enabled:
|
for name in enabled:
|
||||||
try:
|
try:
|
||||||
src = get_source(name)
|
src = get_source(name)
|
||||||
@@ -49,6 +34,31 @@ def fetch_node(state: Dict[str, Any]) -> Dict[str, Any]:
|
|||||||
})
|
})
|
||||||
print(f"[fetch] 数据源 {name} 失败(跳过): {e}")
|
print(f"[fetch] 数据源 {name} 失败(跳过): {e}")
|
||||||
|
|
||||||
|
if rows:
|
||||||
|
rows = validate_rows(rows, "fetch")
|
||||||
|
stats = dict(state.get("stats") or {})
|
||||||
|
stats["fetch"] = {"raw_rows": len(rows), "sources": enabled, "errors": len(errors)}
|
||||||
|
return {"raw_rows": rows, "errors": errors, "stats": stats}
|
||||||
|
|
||||||
|
# 2) 抓取失败/无结果 → 回退采集缓存(filter 上次写入的 collected_keywords.json)
|
||||||
|
use_collected = (config.get("fetch") or {}).get("use_collected", True)
|
||||||
|
if use_collected:
|
||||||
|
try:
|
||||||
|
import json as _json
|
||||||
|
from pathlib import Path as _Path
|
||||||
|
p = _Path(state.get("cache_dir") or state.get("output_dir", "")) / "collected_keywords.json"
|
||||||
|
if p.exists():
|
||||||
|
data = _json.loads(p.read_text(encoding="utf-8"))
|
||||||
|
cached_rows = data.get("keywords") or []
|
||||||
|
if cached_rows:
|
||||||
|
rows = [dict(r) for r in cached_rows] # 已过滤去重的关键词
|
||||||
|
print(f"[fetch] 数据源抓取失败,回退采集缓存 {len(rows)} 条({country})")
|
||||||
|
stats = dict(state.get("stats") or {})
|
||||||
|
stats["fetch"] = {"raw_rows": len(rows), "sources": ["collected_cache"], "errors": len(errors)}
|
||||||
|
return {"raw_rows": rows, "errors": errors, "stats": stats}
|
||||||
|
except Exception as e: # noqa: BLE001
|
||||||
|
print(f"[fetch] 读取采集缓存失败(回退数据源): {e}")
|
||||||
|
|
||||||
rows = validate_rows(rows, "fetch")
|
rows = validate_rows(rows, "fetch")
|
||||||
stats = dict(state.get("stats") or {})
|
stats = dict(state.get("stats") or {})
|
||||||
stats["fetch"] = {"raw_rows": len(rows), "sources": enabled, "errors": len(errors)}
|
stats["fetch"] = {"raw_rows": len(rows), "sources": enabled, "errors": len(errors)}
|
||||||
|
|||||||
Reference in New Issue
Block a user