技术札记

实战多 API 数据采集器:异构来源如何汇成一张增量表

用 HHY 并发采集 OpenAlex、Crossref 和 GitHub,统一字段、隔离分页失败,并按复合键增量合并 CSV。

HOUHUIYANG.COM

扫码继续阅读

正在生成…

实战多 API 数据采集器:异构来源如何汇成一张增量表

houhuiyang.com/zh/notes/building-a-multi-api-data-collector-with-hhy

采集一个 API 不难。真正有工程价值的问题是:三个来源的分页规则、JSON 结构和指标完全不同,如何让它们进入同一条可靠 Flow,并且第二次运行不会把第一次结果简单覆盖?

这个项目同时采集 OpenAlex 学术成果、Crossref 论文元数据与 GitHub 仓库。默认每个来源两页、每页五条,最终统一成七个字段:sourceexternal_idtitleurlrecord_typemetric_namemetric_value

多 API 数据采集器的真实项目目录

先把分页差异变成 Job

OpenAlex 和 GitHub 使用 page,Crossref 使用 offset。jobs.hhy 不试图抹掉这个差异,而是把它变成显式任务:

for page in range(1, config.pages + 1) {
    jobs = append(jobs, { source: "OpenAlex", page: page, offset: 0 })
    jobs = append(jobs, { source: "Crossref", page: page, offset: (page - 1) * config.per_page })
    jobs = append(jobs, { source: "GitHub", page: page, offset: 0 })
}

六个 Job 进入 parallel(3)。每个请求前等待 250 ms,并设置 10 秒 timeout、两次 retry 与 500 ms backoff。这不是严格的全局 rate limiter,但能避免示例在启动时瞬间向公开 API 发出一组突发请求。

GitHub 未认证搜索额度较低。生产环境需要更低频率或通过安全的进程扩展注入认证请求,而不是把 token 写进配置和仓库。

统一 Schema,但不伪装语义相同

三种 API 分别有自己的 request 与 normalize 函数。OpenAlex 的指标是 citations,Crossref 是 references,GitHub 是 stars;它们共享 metric_name/metric_value 形状,却没有被错误地塞进一个含义模糊的 score

OpenAlex → work id / display_name / citations
Crossref → DOI / title / references
GitHub   → full_name / repository / stars
                    ↓
source / external_id / title / url / type / metric_name / metric_value

统一字段是为了让存储、去重和消费简单;保留指标名是为了不牺牲来源语义。

下载成功后解析仍可能失败,因此网络阶段和 normalize 阶段都有独立的 attempt。失败结果记录 source、page 和 error,并写入 failures.json。成功页仍可继续进入合并,但总报告会 ok=false 并返回 1。

增量合并需要稳定身份

程序先尝试读取已有 CSV,再把旧记录与新记录展开到同一条 Stream。复合键使用 source + external_id

[existing, incoming]
    |> stream
    |> flat_map { records -> records |> stream }
    |> group_by { record -> join([record.source, record.external_id], ":") }
    |> map { group -> group.values[length(group.values) - 1] }
    |> sort_by { record -> join([record.source, record.external_id], ":") }
    |> collect

旧记录在前、新记录在后,因此同键时保留最后一条,新数据自然覆盖旧数据。最后按相同复合键稳定排序,避免 API 返回顺序变化制造无意义 diff。

CSV、运行报告和失败列表都原子写入。输出记录保留原始 URL,让每一行都能回到来源核验。

多 API 数据采集器连续两次运行的增量验证

为什么测试要连续运行两次

本地 fixture 第一次返回 12 条跨页输入,去重后得到 9 条唯一记录。第二次运行先读取这 9 条旧记录,再采集、覆盖和合并,最终仍然是 9 条。

sh practical-projects/multi-api-data-collector/self-test.sh

只测试第一次运行,只能证明“可以导出 CSV”;连续两次才验证复合键、覆盖方向、稳定排序和幂等形状。测试服务使用本机随机端口和临时目录,退出后删除,不会把 fixture 数据写进正式 output。

这条采集 Flow 的边界

parallel 在这里优化的是 HTTP 等待,不是把数值计算包装成并发故事。retry 处理短暂故障,不会修复错误 Schema。增量合并保留最新观测,也不等于拥有完整历史。

把边界讲清楚以后,这个项目才真正展示 HHY 的方向:用少量可组合原语处理网络、Stream、错误和原子文件,让可靠性策略直接出现在业务 Flow 里。

参考

返回技术札记