← 代码课堂

第 17 课:数据流水线 —— 采集、去重、可解释打分与预算护栏

难度:★★★(综合课:数据处理 + 工程 + 成本意识)
教材:paper-follower/(定时追踪论文的机器人,172 个测试)
预计时间:讲解 90 分钟 + 练习 50 分钟
学完你能:设计一条"每天自动跑"的数据流水线,并处理重复、误判、成本和漏跑


课前须知:定时任务比看起来难

"每天自动抓一批论文,用 AI 判断哪些值得看"——听起来简单,但真实工程要回答:

这一课逐个看它怎么解决。


第一步:跑起来

cd paper-follower
py -3 -m venv .venv
.venv\Scripts\pip install -r requirements.txt

.venv\Scripts\pytest -q          REM 172 个测试,看全绿
.venv\Scripts\python -m app.cli --help
.venv\Scripts\python -m app.cli health       REM 体检:它到底在正常跑吗

.env 里配好 DEEPSEEK_API_KEY 后,可以试:

.venv\Scripts\python -m app.cli daily --skip-if-done

第二步:模块地图

app/
├── sources/        采集:openalex.py(主)、arxiv.py(补预印本)
├── pipeline.py     流水线:抓取 → 合并 → 打分 → 入库
├── normalize.py    归一:DOI/arXiv/标题指纹、词干、倒排摘要还原
├── store.py        SQLite:papers / summaries / llm_usage / runs 四张表
├── scoring.py      可解释打分(本课重点)
├── summarize.py    三档摘要(skim / brief / deep)
├── llm.py          模型调用 + 记账 + 预算护栏 + 分时定价
├── digest.py       生成 Markdown 日报
├── export_course.py 导出成 course-publisher 的课程源
└── cli.py          命令入口 + 定时守卫 + health

install_task.bat    注册三个 Windows 计划任务(早/晚/开机后)
run_daily.bat       计划任务实际执行的脚本

第三步:逐块精讲

块 1:多源采集与合并(pipeline.py)

思路:OpenAlex 当主力(按主题词多路并发抓),arXiv 补预印本,
同一篇论文可能两边都有——于是要先合并:

两篇记录 → 用同一套字段合并:谁的字段全用谁的,引用数取更大的

对应代码在 pipeline.py 的 collect() / _merge() / _finalize()。
用了 ThreadPoolExecutor 并发抓取——网络 IO 密集任务,并发能明显提速。

设计要点:不同数据源的字段名不一样(有的叫 title,有的叫 name),
normalize.py 负责把它们统一成内部格式。先归一化,再处理,是数据流水线的通用步骤。

块 2:去重 —— 判断"两篇是不是同一篇"(normalize.py 第 217-223 行)

def paper_uid(doi: str = "", arxiv_id: str = "", title_norm: str = "") -> str:
    """内部唯一键:DOI > arXiv id > 标题指纹"""
    if doi:
        return f"doi:{norm_doi(doi)}"
    if arxiv_id:
        return f"arxiv:{norm_arxiv_id(arxiv_id)}"
    return f"title:{title_norm[:120]}"

三级优先级:正式的 DOI 最可靠 → arXiv 编号 → 标题指纹(去掉大小写/标点后)。
但这还不够:预印本常常没有 DOI,正刊有 DOI,两边标题一样——
所以入库时(store.upsert)还要用 title_norm 兜底:
uid 没撞上时,再看标题是否撞上;而且规则是"已有正式 DOI 的记录不被无 DOI 的覆盖"。

通用的去重套路:先定一个稳定主键,再留一条兜底规则。
只靠主键会漏(id 格式不同),只靠兜底会误合(标题相似但不是同一篇)。

顺带一个实用函数(第 204-214 行):

def abstract_from_inverted_index(inv: dict | None) -> str:
    """OpenAlex 的摘要是倒排索引 {词: [位置...]},还原成正常文本"""
    ...
    slots[pos] = word
    return " ".join(w for w in slots if w)

OpenAlex 为了省流量,摘要存成"词 → 出现位置"的倒排索引。
这个函数把它还原成正常句子——处理第三方数据格式,是数据工程的日常。

块 3:可解释打分(scoring.py 第 1-9、177-244 行)

文件开头第一句就是纪律:

核心纪律:每条推荐必须能说出人话理由,写不出理由就不许进 digest。

WEIGHTS = {
    "topic": 0.35, "citation": 0.20, "seed": 0.15,
    "author": 0.15, "repro": 0.10, "novelty": 0.05,
}

def score_paper(paper, profile, seed_ids=None, cross_pairs=None, low_trust_venues=None, ...):
    topic, hit_title, hit_abs = _topic_score(paper, profile)
    citation = _citation_score(paper)
    seed, seed_reasons = _seed_score(paper, profile, seed_ids)
    ...
    parts = {"topic": topic, "citation": citation, "seed": seed, "author": author,
             "repro": repro, "novelty": novelty}
    score = sum(WEIGHTS[k] * v for k, v in parts.items())

    trust, trust_note = venue_trust(paper, low_trust_venues, low_trust_factor)
    score *= trust                      # 无同行评审的平台打折
    ...
    reasons: list[str] = []
    if hit_title: reasons.append(f"标题命中主题词:{hit_title[0]}")
    if hit_abs:   reasons.append(f"摘要命中主题词:{hit_abs[0]}")
    reasons.extend(seed_reasons)
    ...
    return {"score": round(score, 4), "score_detail": {...}, "reasons": reasons, ...}

三个值得抄的设计:

  1. 加权求和:六个维度各有权重,分数=Σ(权重×维度分)。

比"if 一堆条件"更可调、更可解释。

  1. trust 折扣:score *= trust——arXiv/bioRxiv 这类没有同行评审的平台打 0.55 折。

但又不打到 0、也不丢弃(注释说得很清楚:那里也有重要的东西)。

  1. 每个分数都配一句中文理由(reasons)——这是"可解释"的落地方式:

最终日报里,每条推荐后面都跟着"为什么推荐它"。

块 4:两套匹配规则 —— 一个"真实翻车"的教学范本(normalize.py 第 150-201 行)

这是全项目最精彩的一段注释,强烈建议原文读一遍。核心矛盾:

作者的解法:两套规则,用在两个地方。

宽松版 term_hits(第 150-175 行)——用于主题匹配:

"""规则:
- 词组里若有高特异性词(如 bioelectric),命中任一个就算数——
  因为摘要常常只写 bioelectrical 而不写完整的 "bioelectric signaling"。
- 若全是泛词(如 membrane potential / protein structure),则要求全部命中,
  否则单凭 potential 一个词就能把核聚变论文放进来。
这是**召回**用的宽松规则,故意比 term_all_hits 宽。"""

严格版 term_all_hits(第 178-201 行)——用于"交叉监控"(finding 机会窗口):

"""交叉监控用的**严格**判定:词组里每个实词都必须命中。
...
实测把交叉监控切到那条规则后,"ion channel expression × development"
立刻捞进了《阳光暴晒对日灼病的影响》《First person – Fabienne Dreier》。
交叉监控宁可漏也不能错:...混进噪音就等于整个机会窗口失去意义。所以这里要求全中。
严格的是词数,不是形态:signalling 与 signaling、regenerative 与 regeneration
这类形态差异照样认(见 stem / same_stem),因为那是同一个概念。"""

这段代码 + 注释值一节课,它体现了三个工程素养:

  1. 同一个问题在不同场景需要不同的准确率/召回率取舍
  2. 记录"为什么",尤其是"我试过什么、翻过什么车"——这样后人不会把它改回去
  3. 区分"词数严格"和"形态严格":概念相同(signalling/signaling)要认,

泛词堆砌要拒

配合 normalize.py 的 stem / squash_doubles / same_stem(第 77-142 行)实现"词干化":
regenerative ↔ regeneration 能对上,generation ↔ regeneration 不会误合。

块 5:LLM 记账、分时定价、预算护栏(llm.py 第 24-137 行)

# 元 / 百万 tokens。空闲时段是高峰时段的一半。
# 高峰 = 周一至周五 9:00-12:00、14:00-18:00(不含法定节假日),其余全为空闲。
PRICES = {
    "deepseek-flash": {
        "peak": {"in_hit": 0.04, "in_miss": 2.0, "out": 8.0},
        "off":  {"in_hit": 0.02, "in_miss": 1.0, "out": 4.0},
    },
    ...
}

def is_peak(when=None) -> bool:
    now = when or datetime.now()
    if now.weekday() >= 5:      # 周末全天空闲
        return False
    minutes = now.hour * 60 + now.minute
    return (9*60 <= minutes < 12*60) or (14*60 <= minutes < 18*60)

def estimate_cost(model: str, usage: dict) -> float:
    table = PRICES.get(model, FALLBACK_PRICE)
    p = table["peak"] if is_peak() else table["off"]
    hit  = int(usage.get("prompt_cache_hit_tokens") or 0)
    miss = usage.get("prompt_cache_miss_tokens") or max(0, int(usage.get("prompt_tokens") or 0) - hit)
    out  = int(usage.get("completion_tokens") or 0)
    return hit/1e6*p["in_hit"] + int(miss)/1e6*p["in_miss"] + out/1e6*p["out"]

注意三个成本维度:缓存命中的输入最便宜、未命中的输入贵、输出最贵。
所以"关掉思维链"(第 110 行 "thinking": {"type": "disabled"})能实打实省钱。

预算护栏(第 96-103 行):

if check_budget:
    cap = settings.daily_cost_cap_cny
    spent = spent_today()
    if cap > 0 and spent >= cap:
        raise BudgetExceeded(
            f"今日已花 ¥{spent:.3f},达到上限 ¥{cap:.2f}。"
            f"明天再来,或调高 .env 里的 DAILY_COST_CAP_CNY。"
        )

每次调用前先查"今天花了多少",超了就抛异常。
错误信息还告诉用户怎么解决——错误提示要可操作。

重试策略也很讲究(第 127-136 行):

except urllib.error.HTTPError as e:
    ...
    if e.code in (400, 401, 403, 404):   # 请求本身有问题,重试没意义
        raise last_err
    time.sleep(2 * (attempt + 1))        # 其他错误退避重试

分清"该重试的错误"和"重试也没用的错误"(参数错、没权限 → 别浪费时间和钱)。

块 6:定时任务的幂等守卫(cli.py 第 534-553 行)

Windows 计划任务有个坑:错过了的触发时间不会补跑。
作者的对策是注册三个触发点(早 7:30 / 晚 19:30 / 开机后 2 分钟),
谁先跑谁干活,其余秒退——靠 --skip-if-done 和这个守卫:

def already_done_today() -> tuple[bool, str]:
    """今天是否**完整**跑过一轮,返回 (是否完成, 说明)。

    "完整" = 抓取和日报两件都做了。只做了一半(比如抓完中途挂了)不算完成 ——
    否则晚上那次重试就被自己挡掉了,而那正是最需要它的时候。

    判断用的是两处独立证据,不猜:
      - runs 表今天的记录  → 抓取成功了
      - 今天的 digest 文件 → 日报生成成功了
    """
    fetched = _fetched_today()
    digested = _digested_today()
    if fetched and digested:
        return True, "今天已完成(抓取 + 日报都在)"
    ...

两个独立证据(数据库记录 + 文件存在)比"一个标志位"可靠得多:
只有两件都完成才算完成,半途而废的还能被晚场重试救回来。

配套还有 health 命令(第 556-561 行)——理由写得很实在:

定时任务没生效时,系统不会有任何报错——它只是安静地什么都不做。
没有这个命令,只能靠"感觉最近没收到日报"来发现。

块 7:跨项目只读耦合(export_course.py + phone-remote 的桥接)

paper-follower 会把带读导出成 course-publisher 的课程源
(写文件,不 import 对方代码);phone-remote 则以只读模式打开它的
papers.db(唯一允许的写操作是 UPDATE papers SET state=?)。

项目之间通过"数据边界"合作,而不是互相 import——
这样任何一个项目单独删掉/重写,另一个都不会崩。这是架构上的重要选择。


第四步:对照开源,别人怎么做数据流水线

对照 ETL / 数据管道

工业界把这类流程叫 ETL:Extract(抽取)→ Transform(转换)→ Load(加载)。
教材项目恰好是:sources/(抽取)→ normalize + scoring(转换)→ store(加载)。
概念一致,只是换了业务场景。

对照"可解释性"

方案特点
纯规则打分(你的)完全可解释、可调权重、零成本
机器学习排序效果可能更好,但难解释、要训练数据
纯 LLM 打分灵活,但贵、不稳定、难复现

你选了"规则打分 + LLM 只做摘要",在成本、可解释、稳定之间取得平衡——
这是很成熟的工程判断。

对照幂等与可观测

你用一个 951 行的 cli.py 做了个精简版


动手练习(做完才算过关)

练习 1(热身):跑测试 + 体检
pytest -q 看到全绿;然后运行 python -m app.cli health,观察它检查了哪些项。

练习 2(必做):调权重看排序
把 scoring.py 的 WEIGHTS 里 citation 从 0.20 调到 0.50,
跑一次打分(或测试),观察高分论文的排序变化。体会:权重就是产品策略。

练习 3(必做):匹配规则实验
在 config/interests.yaml 的主题词里加一个泛词(比如 development),
用 term_hits 和 term_all_hits 分别在几篇论文上进行判定(写个小脚本调用即可),
对比两者的命中差异。亲眼看看"宽松"和"严格"的区别。

练习 4(必做):触发预算护栏
把 .env 的 DAILY_COST_CAP_CNY 设成很小的值(比如 0.0001),
跑一次需要调用模型的操作,观察 BudgetExceeded 的错误信息。
(记得改回来。)

练习 5(思考题):
为什么 already_done_today 要用"数据库记录 + 文件存在"两处证据,
而不是写一个 done=true 的标志位?

看答案

标志位什么时候写?如果写完标志位后程序崩了怎么办?

练习 6(选做):给日报加字段
在 digest.py 生成的日报里,给每篇论文加上 score_detail 的各维度分数
(你已经拿到了这个字典)。先看它现在的 Markdown 结构,再决定加在哪。


本课小结

术语表

术语人话解释
流水线 / ETL抽取→转换→加载的自动化数据处理流程
归一化(normalize)把不同来源的数据统一成同一格式
倒排索引"词→位置"的存储方式,OpenAlex 用它存摘要
词干(stem)把词还原成词根,如 signalling→signal
召回 / 精确别漏掉 / 别捞错,两者常需权衡
幂等重复执行结果不变
预算护栏花费超限就拦住,防止意外烧钱
退避重试失败后等越来越久再试
可观测性能判断系统是否正常运行(日志/健康检查)
数据边界耦合项目间靠数据文件/库协作,不 import 代码

下节预告

第 18 课:纯函数与脱离平台测试。
教材是 wechat-map-tools(微信小程序地图工具,47 项单测):
坐标怎么在 WGS-84 / GCJ-02 / BD-09 三种系统间转换(含"不动点迭代"这个小技巧)、
距离和面积怎么用球面公式算、以及怎么让逻辑脱离微信环境就能测。