aggregated-job-ingestion-pipeline

聚合职位摄入管线

多个平台上的爬虫抓到教师岗位,把每一条抓取结果丢进一个队列。队列的另一端是一个职责单一、清晰的 worker:把一条原始抓取变成一行规范的岗位记录,并让这行记录立刻对所有读取岗位的一方可见——缓存、搜索索引——而再另起一个队列去做这些后续工作。本节讲的就是这道接缝,以及烤进其中的那些取舍。

从队列到单独一行

摄入 worker 是一个独立进程。它消费一个由 Redis 支撑、名称来自配置的 Bull 队列,只处理一个名为 job-update 的任务,并发有上限(默认 5)。启动时它把这份工作所需的一切接好,此后不再做别的:一个跑在 Postgres 上的 Prisma 仓库、一个 Redis 岗位缓存和一个 Redis 搜索缓存,以及——如果连得上 Meilisearch——一个 Meilisearch 索引器。若 Meilisearch 客户端初始化失败,worker 记下日志、以「索引已禁用」继续跑,而不是拒绝启动。一个 DISABLE_QUEUE_CONSUMER 开关可以让进程保持存活但空转;SIGINT/SIGTERM 触发一次有序退出,关闭队列、断开 Postgres 与 Redis 连接。

每一条排队的载荷都经过 AggregatedJobIngestService.process。这个服务把原始抓取规范化成一份插入形状,确保雇主存在,upsert 单独一行,然后刷新读路径。雇主这一步是有意为之:当带有学校名时,服务同步地确保存在一行雇主档案,并在插入之前把它的 id 盖到岗位上。代码注释记下了原因——它替换掉了旧的「尽力而为的查找 + 异步的学校摄入队列」,于是每条岗位落地时都已链接到它的雇主,不存在那个「首次抓取的岗位悬空未链接」的竞态窗口。

确定性的、人可读的 ID

这行记录的主键不是随机 UUID。当来源给出一个稳定的 token(平台岗位 id,或外部 id)时,id 由它确定性地构造出来:一个平台前缀、token 的一段 slug 化的人可读切片,加上 token 的一小段 SHA-256 摘要,拼起来并截断到 120 字符。同一个「平台 + token」永远算出同一个 id——这正是关键所在。对同一条帖子的再次抓取会算出同一个 id,命中同一行 Postgres 记录,upsert 就地更新它,而不是新建一条重复。去重是 id 本身的性质,不是一次单独的查找。

当来源没有给出稳定 token 时,服务退回到一个以时间为种子的 id——平台 slug、标题 slug,再加一个 36 进制的时间戳。这个 id 可读,但在多次抓取之间并不稳定,于是无 token 的来源用「不去重」换来了「本就没有可用于去重的稳定键」这一事实。这套设计诚实的形状是:只有当来源给出可作键的东西时,去重才被保证。

其余的整理由入口处的规范化一次做完:货币代码统一大写(RMB 折叠为 CNY),薪资数字从嘈杂的字符串里解析出来,位置文本由「城市 / 省 / 国家」拼成,过期时间盖在发布日期之后 60 天,原始抓取则保留在这行记录上,这样规范化不丢任何东西。

就地刷新,不设二级队列

upsert 之后,服务调用一个 notifier——jobUpserted——由它把这行新鲜记录扇出到各个读路径:Redis 岗位缓存、Redis 搜索缓存、Meilisearch 文档。这一切在同一个 worker 里、在与写入同一个工作单元里、内联发生。没有第二个「现在去索引它」的队列,也没有单独一个「现在去预热缓存」的任务。写入与刷新是同一步,于是一条岗位在它那行记录存在的那一刻就可被搜索到。

唯一内联的,是 Discord 通告。它只在 upsert 确实新建了一行时触发(再次抓取的更新不触发),而且以「发了就不管」的方式触发——再次抓取不能反复刷屏,首次摄入也不能被「等一条通知发出去」拖住。

从不出门的富化

同一个 job 服务里的两个辅助服务共享一个立场:富化必须在摄入路径上不做网络往返就能回答。

地理编码完全离线且同步。 它把一份 GeoNames 城市数据集一次性载入内存,通过一条回退阶梯从本地数据解析坐标——先「城市 + 国家」,再单独「城市」,再把城市字段当省来读(爬虫常把省错填进城市字段),再显式的省,最后是「只有国家」时回退到一个代表性城市。旧的 Nominatim(走网络)路径被显式标为弃用;现在那个异步入口只是包住同步查找。位置字符串会先被规范化——去掉变音符号、去掉撇号、转小写——这样查找对爬虫如何拼写一个地名是宽容的。

汇率是被缓存的,不是每条岗位都去取。 汇率从一个公开的货币汇率 CDN 拉取,在 Redis 里缓存七天,上面再叠一层内存副本,于是一次读取先查内存、再查 Redis,只有冷未命中时才去取。刷新最多一天尝试一次,刷新失败只记日志、不抛出——在热路径上,「陈旧但存在」的汇率优先于一次硬失败。

底下的设计

有一个立场贯穿始终:让那一行记录成为唯一真相之源,让它的 id 是确定性的,好让再次抓取收敛而非增殖,并把每一项后续——链接、缓存、索引、通告——都内联在那一次写入之后完成,而不是通过一堆二级队列扇出。 富化不出门,于是热路径永不因第三方而阻塞;唯一被允许异步、尽力而为的,是那条面向人的通告,而下游没有任何东西依赖它。

about this entry

One of sijie's wiki entries. The AI on this site is grounded in the same corpus and answers in sijie's voice, with citations back to entries like this one — answering costs sijie money, so it waits behind a code: enter an access code →

aggregated-job-ingestion-pipeline