调度、队列扇出,与 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 面很小且大多异步——它请求机器去做事、并回报它接受了什么,而不是在请求内部把活干完。