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

scrape-parse-store-pipeline

抓取 → 解析 → 存储的流式管线

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、记为一条失败结果,循环接着跑下一个平台,而不是拖垮整批。

底下的设计

这套管线自始至终守着一个立场:流式而非批量,好让内存保持平坦、让工作可续跑;让数据库的唯一性规则来做去重,好让再次抓取收敛到同一行;把「消失」当作「没有被重新上线」而不是一次显式删除。 重试、逐条职位的错误隔离、带超时的资源关闭、平台之间的隔离,全都服务于同一个目标——一条闪断的职位、一个卡死的浏览器、一个坏掉的平台,都绝不拖垮它周围的整次运行。

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 →

scrape-parse-store-pipeline