抓取 → 解析 → 存储的流式管线
ScraperPipeline 把一个平台从头跑到尾:一个爬虫把职位列表一条一条流式吐出,每一条在下一条到来之前就被解析并写入 Postgres,不再出现的职位则被下线。整套运行的设计目标是:任何时刻内存里只有一条职位的 HTML;一次中途被打断的运行能从断点接着跑。本节讲的就是这套编排,以及烤进其中的那些取舍。
一次只处理一条,是刻意的
爬虫这一侧是一个异步生成器。ScraperStrategy.scrape(browser) 声明为 AsyncGenerator<RawJobListing>,每个平台在内部自行决定怎么喂:Dave's ESL Cafe 先抓一次列表页,然后把每个详情页一条一条地取下来并 yield 出去。基类上的注释把意图说得很直白——为了内存效率,一次只 yield 一条,于是任何时刻内存里只有一条职位的 HTML。
管线用一个 for await 循环消费这个生成器。对每一条被 yield 出来的 RawJobListing,它计数、记下 jobId,然后要么跳过它(续跑模式,见下文),要么交给 processOneJob。因为生产和消费是交错进行的,管线从不把完整的职位列表实体化——两次 yield 之间爬虫被挂起,限流也活在爬虫内部(Dave's ESL 在两次详情页抓取之间按配置的 rateLimit 睡 5 秒)。
Dave's ESL 爬虫本身每一步都很防御:它等 Angular 应用渲染出来,对列表尝试多个候选选择器,只保留匹配详情页 URL 模式的链接,并从该 URL 推导出每条职位的 id。某个详情页加载失败会被记日志并跳过——循环继续处理下一条,而不是让整次运行失败。只有列表页本身失败时爬虫才会重新抛出。
解析并 upsert,带退避
processOneJob 把它的工作包在 retryWithBackoff 里:最多三次尝试,延迟每次翻倍(1 秒,然后 2 秒)。在被重试的块里,它先查这条 (平台, jobId) 是否已有记录,解析原始 HTML——把上一次的职位描述和原始数据传进去,好让解析器复用——再把发布日期和原始 HTML 盖到解析出的职位上并 upsert。upsert 的条数从数据库返回,累加进总数。如果三次尝试都失败,错误被记日志、这次运行的 errorCount 加一;一条坏职位不会中止整次运行。
记录的身份是一条数据库唯一性规则,而不是应用层的记账:Job 表在 (sourcePlatform, jobId) 上唯一,于是对同一条帖子的再次抓取会就地更新同一行。每行还带一个 active 标志、一个 lastSeenAt 时间戳和一个 parsedWithAI 标志;下线逻辑靠的是前两样。
先全部下线,边跑边重新上线
一次全新的运行会做一件看似反直觉的事:在抓取之前,先把该平台当前所有 active 的职位标记为 inactive。然后,随着每条抓到的职位被 upsert,它又变回 active。效果是:这次运行里爬虫没有看到的东西,会被留在 inactive 状态,而无需一趟单独的 diff——upsert 这条流本身就是 diff。管线把 upsert 的条数记为「re-activated(重新上线)」,正是出于这个原因。
到最后,finalizeRun 会跑第二趟更慢的清扫:deactivateStaleJobs 拿到这次运行里见过的每一个 jobId,把那些没被看到、且 60 天以上未被看到的职位下线。下线总数是「开头那次全量下线」与「这趟陈旧清扫」之和。于是有两个时间视界:「这次没见到」(由「先下线再上线」的手法立即处理)和「很久没见到」(结尾处的 60 天规则)。
从未完成的运行续跑
每次运行都作为一行 ScrapeRun 记录被追踪,带状态和计数器。启动时,管线向数据库询问该平台是否有未完成的运行。若有,它就续跑这一次:复用那次运行的 id 而不新建,并载入自那次运行开始以来已处理过的职位集合。当爬虫把这些职位重新 yield 出来时,它们被跳过——计入 parsed 和 upserted 以让总数保持诚实,但不再重新解析或重新写入。
关键在于,续跑路径跳过开头那步「全部标记 inactive」。那次下线在原始运行开始时就已经发生过了;再做一次会错误地把这次续跑还没走到的所有职位打下线。这就是「开头的全量下线」和「续跑分支」互斥的原因——设计把「全部下线」当作「每个逻辑运行一次」的动作,而不是「每个进程一次」的动作。
清理和失败都不是可选项
资源处理是显式的。浏览器每次运行以 headless 方式启动一次,并总在 finally 里关闭,走一个 safeClose 辅助函数:它让关闭动作与一个超时赛跑,若关闭卡住就记日志而非抛出——这样一个卡死的浏览器永远塞不住进程。当运行本身抛出时,handlePipelineError 会先把失败记到该平台的元数据上(状态 failed、一条错误消息)再重新抛出,于是失败在平台记录里可见,而不只是躺在日志里。
runMultiplePipelines 坐在上层:它按顺序跑各平台,并让它们彼此隔离——某个抛出的平台会被 catch、记为一条失败结果,循环接着跑下一个平台,而不是拖垮整批。
底下的设计
这套管线自始至终守着一个立场:流式而非批量,好让内存保持平坦、让工作可续跑;让数据库的唯一性规则来做去重,好让再次抓取收敛到同一行;把「消失」当作「没有被重新上线」而不是一次显式删除。 重试、逐条职位的错误隔离、带超时的资源关闭、平台之间的隔离,全都服务于同一个目标——一条闪断的职位、一个卡死的浏览器、一个坏掉的平台,都绝不拖垮它周围的整次运行。