采集一个 API 不难。真正有工程价值的问题是:三个来源的分页规则、JSON 结构和指标完全不同,如何让它们进入同一条可靠 Flow,并且第二次运行不会把第一次结果简单覆盖?
这个项目同时采集 OpenAlex 学术成果、Crossref 论文元数据与 GitHub 仓库。默认每个来源两页、每页五条,最终统一成七个字段:source、external_id、title、url、record_type、metric_name 和 metric_value。

先把分页差异变成 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,让每一行都能回到来源核验。

为什么测试要连续运行两次
本地 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 里。