2026-09-23·by Sijie Wang#node#project#youteacher#job_scrapers

scheduler-queue-and-api

调度、队列扇出,与 REST API 面

这个节点讲 job_scrapers 服务如何自我驱动:它怎么决定什么时候去抓每个平台、抓到的岗位去了哪里、以及人或别的服务能通过 HTTP让它做什么

它如何启动:index.ts

启动时,若缺少一小组必需的环境变量(数据库地址、AI 解析器的凭据),服务直接拒绝启动。若开启了 schema 自动应用,它先确保数据库 schema,再初始化 Redis 任务队列,然后可选地启动后台清理服务。只有在 Express 服务器开始监听之后,它才启动调度器——而且只有在调度器开关没有被显式关掉时才启动。调度器、队列、清理循环各由一个环境开关独立控制,所以一个部署完全可以只作为纯 API 面运行,不带任何自我驱动。

关闭是对称的。收到 SIGTERM / SIGINT 后,进程依次请求调度器停止、清理服务停止、队列关闭、数据库初始化器关闭——每一步的失败都记日志但不致命——然后才退出。没有任何东西被半排空地丢下。

调度器:SchedulerService

调度器不是一个全局循环。它维护每个平台一个定时器——一张从平台名到 setTimeout 句柄的映射——所以每个平台按自己的时钟被调度。启动时它向 scraper factory 询问支持哪些平台,为每个平台确保数据库里存在一行调度任务,然后为每个平台各装一个定时器。

在这个骨架之上有几处设计选择:

  • 错峰(Staggering)。 各平台不会同时开火。每个平台的初始延迟按其下标乘以一个错峰间隔来偏移(下限一秒,默认相隔十五秒),这样第一轮把负载摊开,而不是同一瞬间锤爆所有来源。
  • 持久化的下次运行时间。 如果一个平台存储的任务已经带有 nextRunAt,定时器就按那个时刻触发,而不是错峰偏移。下次何时的真相源是数据库,不是进程内存。
  • 并发上限。 一个集合记录当前正在运行的平台。启动某平台前,服务检查两件事:同一平台没在跑(在跑就跳过),以及当前并发运行的平台数低于可配置的批次上限(默认三个)。到达上限时,该平台不会被丢弃——而是稍后重新排期。工作被推迟,从不丢失。
  • 数据库级别的锁。 跨重启(或多实例)的协调靠任务行上的一个锁列,而不是进程内的互斥量。每轮开始时释放超过可配置时间窗的陈旧锁;仍然新鲜的锁会让这次运行跳过。正是这一点防止两个 worker 同时抓同一个平台。
  • 到期判断与重试。 一个平台被视为到期,当它从未跑过、或其 nextRunAt 已过、或上次运行失败且重试次数仍低于上限。成功后下次延迟是一个频率间隔之后;失败后同样是一个频率间隔之后(下限抬到一分钟),并记录失败以推进重试计数。

每次运行之后——无论成功、失败、还是跳过——finally 块都会重新读取任务,并按其最新的 nextRunAt 重新排该平台的定时器。这个循环自愈:排期总是重新收敛到数据库所说的样子。

优雅关闭会清掉所有定时器,然后在一个有界的时间窗内(三十秒)等待所有在途的 pipeline 完成,再断开连接。一次慢抓取会被给予完成的机会,而不是被立刻杀掉。

结果去哪:队列,jobQueue.ts

抓取并规范化后的岗位不会直接写给消费方。它们被扇出到一个由 Redis 支撑的 Bull 队列上。这把 scraper 的吞吐和下游岗位服务的吞吐解耦开。

队列刻意是可选的。若没有配置 Redis 地址,队列保持禁用,入队调用只是空操作并给出警告——scraper 照样跑,只是无处发布。队列创建会检查一次就绪状态,失败时把自己拆掉,而不是留下一个半开的句柄。

这里有两处保护很重要:

  • 容量守卫。 入队前,入队路径读取队列当前深度(waiting 加 delayed 加 paused),与一个可配置上限比较(默认一千)。超过上限时,入队被跳过并作为背压上报——队列绝不允许无界增长。
  • 自清理的任务。 每个入队任务都设了 remove-on-complete 和 remove-on-fail,所以完成的工作不会在 Redis 里堆积。

载荷是一个扁平、规范化的形状——标题、招聘方、学校、薪资区间、地点各部分、科目、年级、要求等等——带一个 enqueuedAt 时间戳。它是 scraper 与队列消费方之间的契约。

你能让它做什么:routes.ts

Express 路由在 /api/scraper/v1 下暴露一个小小的带版本的面:

  • POST /scrape —— 触发一个或多个平台的抓取。请求体会被校验;平台名被规范化为其规范形式,任何不支持的名字返回 400 并列出支持哪些。被接受的工作以异步方式派发:至少有一个平台被排上时响应 202,一个都排不上时响应 409(例如请求的平台全都已在运行)。调用方拿到的是每个平台的结果列表,而不是一个被阻塞的连接。
  • GET /platforms —— 支持的平台及其存储的元数据。
  • GET /scrape-runs/:platform —— 单个平台最近的运行历史。
  • GET /health —— 存活状态、版本、构建 commit、启动时间。
  • POST /replay —— 把 scraper 自己数据库里当前有效的岗位重新入队到队列,可选地按指定平台过滤。因为下游服务按(平台、外部 id)这对键做 upsert,replay 是幂等的——在抓取时队列曾禁用、或消费方需要从头重建时,重放都安全。
  • GET /diagnose?url=… —— 一个反爬分类器。它并行发四个探针——一个普通 HTTP 抓取、一个宣称自己是 headless 代理的 HTTP 抓取、一次默认的 headless 浏览器访问、以及一次带隐身的——通过比较哪些组合成功来判定一个站点为什么在拦:IP 级封锁、user-agent 封锁、自动化标志检测、还是更深的指纹识别。它报告的是一个结论,不是一个绕过手段;它是一个分诊工具,用来决定某个来源需要哪种抓取策略。

贯穿线

这个组件里有两条脊柱。时间归数据库所有:调度器只是数据库行之上一层薄薄的、自我重排的定时器层,那些行说明每个平台何时到期、是否被锁。流动归 Redis 所有:结果通过一个有界队列扇出,并配一个 replay 逃生口,这样 scraper 和它的消费方可以各自独立地失败与恢复。HTTP 面很小且大多异步——它请求机器去做事、并回报它接受了什么,而不是在请求内部把活干完。

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 →