跳到主要内容

爬取与异步任务:WebCrawler + 自建队列 NuQ

30 秒导读: 单页抓取(01)解决"把一个 URL 变成 LLM-ready 数据";本章解决"把一个站点变成成千上万条数据"。做法是把爬取拆成一张自我生长的作业图——每抓完一页,就从 HTML 里发现新链接、过滤+去重、再把新页当成新作业入队,直到无新页或触顶。撑起这张图的,是 Firecrawl 自己写的作业队列 NuQ(Postgres + RabbitMQ + Redis)。

本章的所有引用都锚定在 commit f4464e19 上。


1. 这是什么(零基础也能懂)

1.1 一句话定义

爬取(crawl)= 从一个起始 URL 出发,顺着页面里的链接一层层往外抓,把整个站点(或其中一个子路径)都变成结构化数据。

1.2 它解决什么问题

假设你要给一个 RAG 系统灌入 docs.example.com全部文档页。你不想手写几千个 URL,你只想说一句"把这个文档站爬下来"。爬取就是干这个的:

  • 你给一个入口 https://docs.example.com;
  • Firecrawl 自动找到站点地图(sitemap)、顺着页内 <a> 链接发现更多页;
  • 每个页面都走一遍单页抓取内核,产出 markdown;
  • 最后把所有页面聚合成一个结果集。

1.3 为什么这件事"难"

一次爬取不是"跑一个函数",而是一个可能持续几分钟、涉及上万次抓取的长任务。它天然带来四个规模问题:

问题白话Firecrawl 的应对(本章主角)
发现我怎么知道站点有哪些页?sitemap + 页内链接提取(WebCrawler)
过滤哪些链接该爬、哪些该丢?filterLinks / filterURL(深度、正则、robots、外链…)
去重同一个页别抓两遍Redis visited 集合 + URL 规范化/排列
调度上万作业怎么排队、限流、不丢自建队列 NuQ

1.4 用起来什么样

对使用者,爬取是异步的:提交后立刻拿到一个 crawl id,再去轮询状态。

# 1) 提交爬取,立刻返回一个 id(不阻塞)
curl -X POST https://api.firecrawl.dev/v2/crawl \
-H "Authorization: Bearer $KEY" -H "Content-Type: application/json" \
-d '{"url":"https://docs.example.com","limit":200,"scrapeOptions":{"formats":["markdown"]}}'
# => { "success": true, "id": "0190...", "url": ".../v2/crawl/0190..." }

# 2) 稍后用 id 轮询,拿到已完成的页面
curl https://api.firecrawl.dev/v2/crawl/0190... -H "Authorization: Bearer $KEY"

这跟单页 /v2/scrape同步返回是根本不同的执行模型——本章 §6 会讲清这条分水岭。

1.5 一句话直觉

把爬取想象成一场"链式反应": 起始页是第一颗中子,它撞出若干新链接(新中子),每个新链接又抓出更多链接……NuQ 队列就是这个反应堆的"控制棒 + 冷却系统",负责让反应既能扩散、又不失控(限流、去重、防丢)。


2. 顶层全景(它大概怎么转)

2.1 两个大块

爬取子系统由两层构成,别混淆:

  • WebCrawler(算法层)——纯粹的"链接学":给一堆 URL,判断哪些该爬、把 HTML 里的链接抽出来、读 robots.txt 和 sitemap。它不碰队列、不抓页面,只做判断。源码:apps/api/src/scraper/WebScraper/crawler.ts
  • NuQ + worker(调度层)——把"抓某个 URL"变成一个作业,排队、限流、执行、去重、收尾。源码集中在 apps/api/src/services/worker/apps/api/src/lib/crawl-redis.ts

2.2 部件一句话职责

部件干什么在哪个文件
WebCrawler链接过滤 / robots / sitemap / 抽链scraper/WebScraper/crawler.ts
crawlController接收 /v2/crawl 请求,建 crawl、发首个 kickoff 作业controllers/v2/crawl.ts
crawl-redis.ts爬取状态:去重集合、作业索引、完成判定lib/crawl-redis.ts
scrape-worker.ts作业执行体:kickoff 扇出 + 抓页后递归扇出services/worker/scrape-worker.ts
NuQ自建作业队列(PG+RabbitMQ+Redis)services/worker/nuq.ts
runNuqWorkerworker 主循环:取作业→执行→标记完成services/worker/nuq-worker-runner.ts
crawl-logic.tsfinishCrawlSuper:收尾聚合 + 完成 webhookservices/worker/crawl-logic.ts
team-semaphore / concurrency-limit团队级并发限流services/worker/team-semaphore.tslib/concurrency-limit.ts

2.3 主线走一遍(高层,不进代码)

一次爬取的生命周期,从提交到收尾:

POST /v2/crawl


[crawlController] 建 StoredCrawl、取 robots.txt、建作业组(group)
│ 入队一个 "kickoff" 作业

[processKickoffJob] 首轮扇出:
│ ├─ 锁定起始 URL,入队它的 single_urls 作业
│ ├─ 探测 sitemap,把 sitemap 里的 URL 也入队
│ └─ 查索引(index)补充已知 URL

[processJob] (每个页面一个) 抓页 → 从 HTML 抽链 → filterLinks 过滤
│ 对每个新链接:lockURL 去重成功者才 addScrapeJob(递归扇出)
│ 自己抓完 → addCrawlJobDone
▼ …这一步不断自我复制,直到没有新链接或触顶 limit…


[finishCrawlSuper] 当"完成数 == 总作业数"→ 聚合、写 DB、发 crawl.completed webhook

记住这张图的三个字:扇出(fan-out)。 爬取没有一个"中央循环遍历所有页";而是每个抓完的页面自己负责发现并入队下一批页面,作业图自己长大。这是理解本章一切的钥匙。


3. 核心原理之一:WebCrawler —— 链接的"守门人"

本节讲:给定一批候选链接,WebCrawler 如何决定"哪些放行、哪些拦下,以及为什么"。

3.1 它要解决的小问题

从一个页面的 HTML 里,<a href> 可能有几百个:有站内的、站外的、图片、mailto:、锚点、社媒分享按钮、超出深度的深层页……绝大多数不该爬。守门人的职责就是把这堆候选压缩成"该爬的那几个"。

3.2 拒绝的理由被显式编码

Firecrawl 没有让过滤逻辑"悄悄 return false",而是给每一种拒绝配了一条面向用户的解释,集中在 DenialReason 枚举里:

crawler.ts:29-41 定义了全部拒绝类型,含义如下:

DenialReason什么时候拒绝
DEPTH_LIMITURL 路径段数 > maxDepth
EXCLUDE_PATTERN命中 excludePaths 里某条正则
INCLUDE_PATTERN指定了 includePaths 却一条都不匹配
ROBOTS_TXT被站点 robots.txt 禁止
FILE_TYPE指向图片/视频/字体/压缩包等非文档扩展名
BACKWARD_CRAWLING跳出了起始 URL 的路径层级,且未开 allowBackwardCrawling
SOCIAL_MEDIA社媒链接或 mailto:
EXTERNAL_LINK跨域,且未开 allowExternalLinks
SECTION_LINK只是 #锚点,视作同页去重
NON_WEB_PROTOCOLmailto:/tel:/ftp:/file: 等非 HTTP 协议

这份"带理由的拒绝"最终会出现在 crawl 的错误报告里,用户能知道为什么某个页没被爬

批量过滤走 filterLinks(crawler.ts:157,签名 filterLinks(sitemapLinks, limit, maxDepth, fromMap, skipRobots, ignoreDiscoveryDepth))。它的核心逻辑其实下沉到了 Rust:

真实实现在 crawler.ts:190,调用从 @mendable/firecrawl-rs 导入的 filterLinks(crawler.ts:18),把 excludes/includes/robots/allowBackwardCrawling 等一股脑传进去,让 Rust 侧做正则和 URL 判断,再把 Rust 返回的原始 denialReasons(如 "DEPTH_LIMIT")在 JS 侧翻译成上表那种带上下文的长句(crawler.ts:218-266,例如 DEPTH_LIMIT 会填入实际 depth 和 maxDepth)。

关键细节:有一条纯 JS 的回退路径。 如果 Rust 调用抛错,代码 catch 后落到 crawler.ts:289 起的手写 JS 过滤(逐条 new URL + 正则 + isRobotsAllowed + isFile)。这是 Firecrawl 反复出现的模式:Rust 快路径 + JS 慢回退,保证一条路挂了功能不崩。

filterURL(单条版,crawler.ts:667)同理,直接调 Rust 的 filterUrl,用于抓页后逐个判断新发现的链接。

3.4 文件类型判断:isFile

哪些扩展名算"文件、不爬",硬编码在 isFile(crawler.ts:770)的 fileExtensions 列表里(.png .jpg .css .js .zip .mp4 .woff …)。

注意几个被注释掉的扩展名——.pdf.docx.xml 没在拦截列表里(crawler.ts:781/789/791),因为 Firecrawl 爬 PDF/Word/XML 当文档内容。这是个容易忽略的设计点。

3.5 robots.txt:抓取、导入、crawl-delay

守门人尊重 robots.txt,分三步:

  1. getRobotsTxt(crawler.ts:455)通过 fetchRobotsTxt${baseUrl}/robots.txt 的内容;
  2. 导入 importRobotsTxt(crawler.ts:505)用 robots-parser 解析,并读出 crawl-delay——它对多个 UA 名做兜底:getCrawlDelay(this.robotsUserAgent) 失败就试 "FireCrawlAgent" / "FirecrawlAgent"(crawler.ts:510-514);同时提取 robots 里声明的 <sitemap>(crawler.ts:516);
  3. 判定 isRobotsAllowed(crawler.ts:757)对单个 URL 用 isUrlAllowedByRobots 判断;若 ignoreRobotsTxt 为真则直接放行。

3.6 sitemap:比爬链接更高效的发现方式

顺着页面 <a> 一层层爬很慢;如果站点有 sitemap,一次就能拿到几千个 URL。tryGetSitemap(crawler.ts:531)就是干这个:

它会尝试多个 sitemap 位置(tryFetchSitemapLinks,crawler.ts:813),按优先级:

起始 URL 自身(若以 .xml 结尾则直接用,否则拼 /sitemap.xml)
├─ robots.txt 里声明的所有 sitemap
├─ 若是子域名 → 再试主域名的 /sitemap.xml(只保留回指子域的 URL)
└─ 兜底:baseUrl + /sitemap.xml

有个防爆炸的护栏:SITEMAP_LIMIT = 25(crawler.ts:20)——一次爬取最多命中 25 个 sitemap 文件,sitemapsHit 集合超过就停(crawler.ts:991)。

sitemap 的解析WebScraper/sitemap.tsgetLinksFromSitemap(sitemap.ts:28):同样是 Rust 快路径 + JS 回退——先用 @mendable/firecrawl-rsprocessSitemap(sitemap.ts:175),失败退到 parseSitemapXml,再失败退到 xml2jsparseStringPromise(sitemap.ts:198)。sitemap 索引(套娃 sitemap)会递归展开(instruction.action === "recurse",sitemap.ts:289)。注意 sitemap 内容本身是用抓取内核 scrapeURL 去下载的(sitemap.ts:105)——爬虫复用了单页抓取的引擎与容错。

3.7 抽链:从 HTML 到候选 URL

抓完一个页面后,要从它的 HTML 里挖出所有链接。extractLinksFromHTML(crawler.ts:728)又是双路:

  • Rust 路 extractLinksFromHTMLRust(crawler.ts:681):调 extractLinks(Rust),再逐个过 filterURL;
  • cheerio 回退 extractLinksFromHTMLCheerio(crawler.ts:693):Rust 挂了就用 cheerio 遍历 $("a"),顺带还会解析 <iframe>data:text/html 内联的 HTML(crawler.ts:712-723)。

4. 核心原理之二:爬取任务的生命周期

本节把 §2.3 那张图落到代码,追一次爬取从提交到收尾。

4.1 提交:crawlController

controllers/v2/crawl.tscrawlController(crawl.ts:30)干的事,按顺序:

  1. 校验请求、检查权限、(可选)用 LLM 从自然语言 prompt 生成爬取参数(crawl.ts:93-132);
  2. limit 压到不超过剩余额度(crawl.ts:175);
  3. 组装 StoredCrawl 对象 sc(crawl.ts:185),这是整个爬取的"配置 + 上下文"快照;
  4. crawlToCrawler(id, sc, flags) 造一个 WebCrawler(crawl.ts:209),并先抓一次 robots.txt 存进 sc.robots(crawl.ts:212);
  5. 解析队列后端、crawlGroup.addGroup(...) 建作业组(crawl.ts:223-233)——一次爬取的所有子作业都归在这个 group 下;
  6. saveCrawl(id, sc) 存 Redis(crawl.ts:235)、markCrawlActive 标记活跃(crawl.ts:237);
  7. 入队一个 mode: "kickoff" 作业(_addScrapeJobToBullMQ,crawl.ts:239),然后立刻给客户端返回 crawl id。

到此 HTTP 请求就返回了——真正的爬取在后台 worker 里进行。这就是"异步"的含义。

术语澄清:函数名叫 _addScrapeJobToBullMQ,但在此 commit 上底层队列已换成自研的 NuQ(见 §5);BullMQ 只是历史遗留的命名。

4.2 首轮扇出:processKickoffJob

kickoff 作业被 worker 取到后走 processKickoffJob(scrape-worker.ts:993)。它是爬取的"点火器",负责铺开第一批要抓的页:

  • 锁 + 入队起始 URL:lockURL(scrape-worker.ts:1013)去重,再 addScrapeJob 入队起始页的 single_urls 作业,并标 isCrawlSourceScrape: true(scrape-worker.ts:1031,标记它是爬取的"源头页");
  • 探测 sitemap:若未禁用,拼出多个候选 sitemap 地址(起始页、/sitemap.xml、根域 sitemap 等),对每个发一个 kickoff_sitemap 作业(addKickoffSitemapJob,scrape-worker.ts:930 / 循环在 1085);
  • 查索引补页:kickoffGetIndexLinks(scrape-worker.ts:900)从 Firecrawl 自己的 URL 索引里捞该站点已知的 URL,批量锁定并入队(scrape-worker.ts:1093-1150);
  • 最后 finishCrawlKickoff(scrape-worker.ts:1154)标记"点火阶段结束"——这是完成判定的一个必要条件(见 §4.5)。

kickoff_sitemap 作业单独由 processKickoffSitemapJob(scrape-worker.ts:1170)处理:用 scrapeSitemap 拉 sitemap → filterLinks 过滤 → 锁定 → 批量入队 single_urls,套娃 sitemap 再递归发 kickoff_sitemap

4.3 递归扇出:processJob 抓完一页后做什么

普通页面作业走 processJob(scrape-worker.ts:218)。抓取本身委托给 startWebScraperPipeline(即 01 的内核,scrape-worker.ts:263)。抓完之后,才是爬取的精华——它就地发现并入队下一批页面(scrape-worker.ts:435-524):

# 示意,非源码:processJob 抓完一页后的扇出逻辑
crawler.setBaseUrl(该页最终 URL) # 处理重定向后的真实基准
links = crawler.filterLinks( # 抽链 + 过滤,一步到位
crawler.extractLinksFromHTML(rawHtml, url))
for link in links:
if lockURL(crawl_id, sc, link): # 去重:抢锁成功才继续
jobId = uuid7()
addScrapeJob({ url: link, mode: "single_urls",
crawlerOptions: { currentDiscoveryDepth: depth+1 } },
jobId, priority) # 入队新页作业
addCrawlJob(crawl_id, jobId) # 登记进 crawl 的作业索引
# 重点看:每个成功抢到锁的新链接,都变成一个新作业 → 图自我生长

真实代码里,lockURL 成功后才 addScrapeJob + addCrawlJob(scrape-worker.ts:461-511),currentDiscoveryDepth 每扇出一层加 1(scrape-worker.ts:487)。抓完自己后,addCrawlJobDone(scrape-worker.ts:636)把本作业记入"完成集合"。

还有个重定向去重的巧处(scrape-worker.ts:410-432):如果页面 A 重定向到 B,会把 A 的所有 URL 排列(generateURLPermutations)塞进 visited,防止 B 又被当新页重复抓;抢不到锁就抛 RacedRedirectError 悄悄放弃。

4.4 去重与状态:crawl-redis.ts

爬取的所有"记忆"都放在 Redis,lib/crawl-redis.ts 是唯一的门面。核心几张 Redis 结构:

Redis key类型作用相关函数
crawl:<id>string(JSON)StoredCrawl 配置快照saveCrawl/getCrawl (:36/:80)
crawl:<id>:visitedset已见过的 URL(去重)lockURL (:449)
crawl:<id>:visited_uniqueset唯一已锁 URL,用于 limit 计数lockURL (:488)
crawl:<id>:jobsset本爬取的全部作业 idaddCrawlJob/getCrawlJobs (:114/:336)
crawl:<id>:jobs_doneset已完成作业 idaddCrawlJobDone (:163)
crawl:<id>:robots_blockedset被 robots 拦的 URLrecordRobotsBlocked (:68)

去重的核心是 lockURL(crawl-redis.ts:449): 它对 URL 先规范化(normalizeURL,crawl-redis.ts:344,可选去掉 query、统一 hash),再 SADDvisited;SADD 返回 1(新加入)才算"抢锁成功",这个作业才有资格入队。这是天然的分布式互斥——多个 worker 并发发现同一个 URL,只有一个能 SADD 到 1,其余得到 0、直接跳过。

lockURL 开头还查 limit(crawl-redis.ts:463-471):visited_unique 的基数 ≥ limit 就直接返回 false,让整张作业图停止生长。

URL 排列 generateURLPermutations(crawl-redis.ts:373) 是个精巧的去重放大器:同一个逻辑页可能有 http/httpswww/非www/、/index.html、/index.php 等写法。当开启 deduplicateSimilarURLs 时,lockURL 用规范排列的第一个变体做键(crawl-redis.ts:478),让这些等价 URL 折叠成一个。函数顶部那段注释(crawl-redis.ts:360-372)还写明了它必须满足的三条不变式,并说明由 permu-refactor.test.ts 证明。

4.5 收尾:什么时候算"爬完了"

判定条件在 isCrawlFinished(crawl-redis.ts:259),需要两件事同时成立:

爬取完成 ⟺ jobs_done 的数量 == jobs 的数量 (所有已知作业都跑完了)
AND kickoff 阶段已结束 + 所有 sitemap 作业已完成

第二个条件由 isCrawlKickoffFinished(crawl-redis.ts:271)保证——光是作业数相等还不够,必须确认"点火阶段没有还在扇出新作业",否则会在图还在生长时误判完成。

真正的收尾动作在 crawl-logic.tsfinishCrawlSuper(crawl-logic.ts:15),由 crawlFinishedQueue 触发(见 §5.5)。它 finishCrawl(crawl-redis.ts:306,标 finish、从活跃集合移除、删 visited 省内存),然后按 v1/v0 分支写 crawl 汇总日志、算总额度、发 crawl.completed / batch_scrape.completed webhook(crawl-logic.ts:128-211)。

有个"零数据保留(ZDR)"的健壮性细节:FDB 后端会在作业完成时抹掉其输入数据,所以收尾时 job.data 可能是 null;finishCrawlSuper 会退回到 StoredCrawl 上持久化的 v1/webhook/requestId 等字段(crawl-logic.ts:40-46)。


5. 核心原理之三:自建队列 NuQ

本节讲全书最"重"的一块:Firecrawl 为什么不用现成队列,而是自己写了一个跨 Postgres + RabbitMQ + Redis 的 NuQ(services/worker/nuq.ts)。

5.1 为什么不用纯 BullMQ

BullMQ(基于 Redis)是 Node 生态常见的队列,Firecrawl 早期也用它(遗留命名到处是 "BullMQ")。但爬取的规模暴露了纯 Redis 队列的短板:

  • 作业量巨大且需持久:一次大爬取几万作业,Redis 内存吃紧;作业状态(结果、失败原因)更适合放关系库持久化查询;
  • 要按团队/爬取做复杂并发控制:见 §5.6,这类"多租户公平调度"用 SQL + Redis 脚本更好表达;
  • 既要吞吐又要低延迟唤醒:大批作业调度用 Postgres,而"某作业完成了、去唤醒等它的人"这种实时信号用 RabbitMQ / PG NOTIFY。

于是 NuQ 把三者组合:Postgres 是事实源(作业表 + 状态机),RabbitMQ 做预取和完成通知,Redis 做并发信号量。

nuq.ts:13 建了一个 pg.Pool(可接 pgbouncer),nuq.ts:6 引入 amqplib,监听/发送分别用 RabbitMQ 或 PG 的 LISTEN/NOTIFY(nuq.ts:123 startListener)。

5.2 作业与状态机

一个 NuQ 作业的形状 NuQJob(nuq.ts:28),状态取值 NuQJobStatus(nuq.ts:22):

queued ──(worker 取走)──▶ active ──(成功)──▶ completed
▲ │
│ └─(失败)──▶ failed
backlog ──(并发放开时提升)──▶ queued
  • queued:等待被取;
  • active:已被某 worker 锁定、正在跑;
  • completed / failed:终态,带 returnvaluefailedReason;
  • backlog:因团队并发上限被压住,存在单独的 <queue>_backlog 表,等有空位再提升(§5.6)。

5.3 入队与出队:靠 SQL 行锁做原子领取

入队就是一条 INSERT(addJob,nuq.ts:808;批量 addJobs 分批 1000 条控制参数量,nuq.ts:930)。

出队是 NuQ 最关键的一句 SQL,getJobToProcess(nuq.ts:1290)。剥掉 RabbitMQ 快路径后,PG 版本是:

WITH next AS (
SELECT ... FROM nuq.queue_scrape
WHERE status = 'queued'
ORDER BY priority ASC, created_at ASC
FOR UPDATE SKIP LOCKED -- 关键:锁住这行,别的事务跳过它
LIMIT 1
)
UPDATE nuq.queue_scrape q
SET status = 'active', lock = gen_random_uuid(), locked_at = now()
FROM next WHERE q.id = next.id
RETURNING ...;

FOR UPDATE SKIP LOCKED 是整个并发领取的命门: 多个 worker 同时执行这句,每个都会锁住并领走不同的一行,谁都不会拿到同一个作业,也不会互相阻塞。这等价于用 Postgres 实现了一个高并发、无重复派发的队列——不需要外部锁。领取时顺手写入一个随机 lock uuid,后续续期/完成都要凭这个 lock 值。

RabbitMQ 存在时还有一层预取:prefetchJobs(nuq.ts:1251)一次用 SKIP LOCKED 领 500 条、推进 RabbitMQ 的 .prefetch 队列,worker 直接从 RabbitMQ get(nuq.ts:1298),把"取作业"的延迟从一次 DB 查询降到一次 MQ 拉取;RabbitMQ 挂了则回退到直接查 PG(nuq.ts:1307)。

5.4 锁续期与完成通知

  • 续期:作业跑得久,得定期证明"我还活着"。renewLock(nuq.ts:1341)UPDATE ... SET locked_at = now() WHERE id AND lock AND status='active'——只有持有正确 lock 才能续。若 worker 崩了不再续期,locked_at 变旧,被回收器当"死作业"重派(这也是 §5.3 里 lock 存在的意义)。
  • 完成/失败:jobFinish(nuq.ts:1366)/ jobFail(nuq.ts:1423)把状态置终态、写 returnvalue/failedreason,并发出通知:有 RabbitMQ 就 sendJobEnd 往该作业的监听通道发一条(nuq.ts:1393),否则用 pg_notify('<queue>', '<id>|completed')(nuq.ts:1390)。
  • 等待结果:同步调用方用 waitForJob(nuq.ts:1141)——"listen 模式"下挂个监听器等通知,"poll 模式"下每 500ms 查一次(nuq.ts:1230)。收到通知后再查一次 DB 拿 returnvalue

5.5 作业组(group)= 一次爬取

一次爬取的几万个子作业,靠 group_id 归属到一个 NuQJobGroup(nuq.ts:1563)。crawlControllercrawlGroup.addGroup(id, teamId, ttl, ...) 就是给这个 crawl 建组(crawl.ts:224)。组的状态 active/completed/cancelled(nuq.ts:1552)。三个队列实例在文件末尾定义(nuq.ts:1734-1739):

  • scrapeQueue(nuq.queue_scrape,带 backlog)——承载所有 single_urls/kickoff 作业;
  • crawlFinishedQueue(nuq.queue_crawl_finished)——收尾信号队列;
  • crawlGroup(nuq.group_crawl)——爬取分组。

当一个组的所有成员作业跑完,系统会入队一个 crawlFinishedQueue 作业;queue-worker.ts 的循环取到它(queue-worker.ts:332)后调 finishCrawlSuper(queue-worker.ts:235)完成 §4.5 的收尾。

5.6 worker 主循环与团队并发

worker 循环 runNuqWorker(nuq-worker-runner.ts:49)是标准的"取—做—标记"循环(nuq-worker-runner.ts:120-199):

while 未关机:
job = queue.getJobToProcess() # 原子领取(§5.3)
if job 为空: 退避 sleep(1.5s→翻倍到10s上限); continue
每 15s: renewLock(job) # 后台续期(§5.4)
result = processJobInternal(job) # 执行(→ processKickoffJob/processJob)
成功 → jobFinish(result) ; 失败 → jobFail(err)

作业分发在 processJobInternalprocessJobWithTracing(scrape-worker.ts:1293/1323)里按 mode 路由到 processKickoffJob / processKickoffSitemapJob / processJob

团队并发限流分两层:

  1. 入队时的准入(lib/concurrency-limit.ts + queue-jobs.ts):addScrapeJob(queue-jobs.ts:500)→ addScrapeJobRaw(queue-jobs.ts:348)先看团队当前活跃作业数是否到 maxConcurrency,到了就不进主队列,而是压进 concurrency 队列(backlog)(_addScrapeJobToConcurrencyQueue,queue-jobs.ts:65,写 NuQ backlog 表 + Redis ZSET)。爬取还叠加一层"每爬取并发上限"(maxConcurrency 或有 delay 时强制为 1,queue-jobs.ts:377-399)。
  2. 完成时的放行(concurrentJobDone,concurrency-limit.ts:291):一个作业跑完,腾出一个槽,就用 getNextConcurrentJob(concurrency-limit.ts:199)从 backlog 里捞下一个提升为 queued。捞取用 ZPOPMIN 原子弹出最小分数成员(concurrency-limit.ts:216),保证多个 worker 不会捞到同一个;捞出来若因"爬取级并发"仍不能跑,就先搁一边、最后再塞回(crawlBlocked,concurrency-limit.ts:259/276)。

同步 scrape 的并发走另一套:team-semaphore.tswithSemaphore(team-semaphore.ts:281)用 Redis Lua 脚本实现带 TTL 的信号量,acquireBlocking(team-semaphore.ts:73)自旋+指数退避+抖动地抢槽,抢到后起心跳线程续租(startHeartbeat,team-semaphore.ts:158)。同步 scrape 占的槽也会"镜像"进异步侧的并发计数,让两条路看到彼此的真实负载(team-semaphore.ts:219-260 的注释与 mirrorSlotAcquire)。


6. 精华:同步 scrape vs 异步 crawl —— 一条执行分水岭

这是本章最该带走的一点:同一个 processJobInternal,在两种模式下的"入口"完全不同。

6.1 异步(crawl / batch)

crawl/batch 走完整的队列:控制器只入队,worker 循环把作业取出来在独立进程里跑。好处是能限流、能重试、能持久化、能横向扩 worker;代价是有队列往返延迟,结果得轮询。

6.2 同步(单页 scrape)

单页 /v2/scrape 不能让用户轮询——它要当场返回结果。所以 scrapeController(controllers/v2/scrape.ts:38)不入队,而是:

  1. teamConcurrencySemaphore.withSemaphore 抢一个并发槽(scrape.ts:227);
  2. API 进程里就地构造一个内存 NuQJob(scrape.ts:258),关键是打了 skipNuq: true(scrape.ts:288);
  3. 直接 await processJobInternal(job)(scrape.ts:302)——同一个执行体,但没走 NuQ 的入队/领取/通知,而是就地同步跑完拿到 doc。

skipNuq: true 会让 processJobWithTracing 跳过并发槽的镜像维护、失败时直接抛错而非序列化回队列(scrape-worker.ts:1401/1422)。

6.3 一句话对比

维度同步 scrape异步 crawl / batch
入口控制器就地 processJobInternal入队 → worker 循环取出
进程API 进程内独立 worker 进程
返回当场返回文档立刻返回 id,轮询取结果
队列绕过 NuQ(skipNuq)走完整 NuQ + backlog + group
并发信号量(team-semaphore)准入/放行两段限流(concurrency-limit)
扇出无(就一个 URL)有(每页递归发现新页)

为什么这么设计: 单页要低延迟、无需持久,就地跑最省;多页要规模、容错、限流、可观测,必须队列化。两者共享同一个作业执行体 processJobInternal,只是"怎么把作业喂进去"不同。


7. 边界与局限(诚实)

  • sitemap 上限硬编码 25(crawler.ts:20)、每爬取最多 20 个 kickoff sitemap 作业(scrape-worker.ts:937TEMP: max 20,注释自认临时):超大站点可能漏发现部分 sitemap。
  • 默认不向上爬:allowBackwardCrawling 关时,起始 URL 路径之外的页会被 BACKWARD_CRAWLING 拦掉——想爬整站得显式开 crawlEntireDomain/allowBackwardCrawling。这常让用户困惑"为什么只爬到几页"。
  • 去重是"尽力"而非绝对:generateURLPermutations 的注释坦承第 3 条不变式(不同 URL 的排列不重叠)"无法证明,超出爬虫范围"(crawl-redis.ts:369-371)——理论上存在极端 URL 折叠误判。
  • 完成判定依赖 Redis 集合基数相等:若某作业既没进 jobs_done 也没被清掉(如进程异常),isCrawlFinished 可能卡住,需靠 TTL(多数 key 24h 过期)和回收器兜底。
  • 三系统耦合的运维成本:NuQ 同时依赖 Postgres、RabbitMQ、Redis(还有 FDB 后端),任一抖动都需回退路径——代码里到处是 catch → fallback,复杂度换来的是可用性。

8. 横向对比

  • 01 抓取内核:本章是"宏观调度",01 是"微观抓取"。爬取的每个 single_urls 作业最终都调用 01 的 scrapeURL;sitemap 内容也用 01 的内核下载。
  • 05 上层能力:map(只发现 URL 不抓)复用了本章的 sitemap + 抽链;crawlmap 的"发现"加上"抓取 + 递归扇出"。
  • 横向(同 shelf,RAG/检索类):很多爬取框架用现成 Celery/BullMQ;Firecrawl 自研 NuQ 换取多租户公平并发 + 关系库持久化,是"规模优先"的取舍。

9. 代码地图(导航索引)

主题文件路径符号名
链接过滤主入口(Rust+JS 回退)apps/api/src/scraper/WebScraper/crawler.tsWebCrawler.filterLinks
单条链接过滤apps/api/src/scraper/WebScraper/crawler.tsWebCrawler.filterURL
拒绝理由枚举apps/api/src/scraper/WebScraper/crawler.tsDenialReason
文件类型判断apps/api/src/scraper/WebScraper/crawler.tsWebCrawler.isFile
robots.txt 抓取/导入/delayapps/api/src/scraper/WebScraper/crawler.tsgetRobotsTxt / importRobotsTxt
sitemap 探测apps/api/src/scraper/WebScraper/crawler.tstryGetSitemap / tryFetchSitemapLinks
HTML 抽链(Rust/cheerio)apps/api/src/scraper/WebScraper/crawler.tsextractLinksFromHTML
sitemap 解析apps/api/src/scraper/WebScraper/sitemap.tsgetLinksFromSitemap
爬取提交控制器apps/api/src/controllers/v2/crawl.tscrawlController
首轮扇出apps/api/src/services/worker/scrape-worker.tsprocessKickoffJob
sitemap 扇出apps/api/src/services/worker/scrape-worker.tsprocessKickoffSitemapJob
抓页后递归扇出apps/api/src/services/worker/scrape-worker.tsprocessJob
作业分发/skipNuq 分支apps/api/src/services/worker/scrape-worker.tsprocessJobInternal / processJobWithTracing
收尾聚合 + webhookapps/api/src/services/worker/crawl-logic.tsfinishCrawlSuper
URL 去重锁apps/api/src/lib/crawl-redis.tslockURL
URL 等价排列apps/api/src/lib/crawl-redis.tsgenerateURLPermutations
完成判定apps/api/src/lib/crawl-redis.tsisCrawlFinished / isCrawlKickoffFinished
NuQ 队列类apps/api/src/services/worker/nuq.tsNuQ
原子领取作业apps/api/src/services/worker/nuq.tsgetJobToProcess / prefetchJobs
锁续期/完成/失败apps/api/src/services/worker/nuq.tsrenewLock / jobFinish / jobFail
等待作业结果apps/api/src/services/worker/nuq.tswaitForJob
作业组apps/api/src/services/worker/nuq.tsNuQJobGroup / crawlGroup
worker 主循环apps/api/src/services/worker/nuq-worker-runner.tsrunNuqWorker
团队并发信号量(同步)apps/api/src/services/worker/team-semaphore.tswithSemaphore
并发准入/放行(异步)apps/api/src/lib/concurrency-limit.tsgetNextConcurrentJob / concurrentJobDone
入队 + backlogapps/api/src/services/queue-jobs.tsaddScrapeJob / addScrapeJobRaw
同步 scrape 就地执行apps/api/src/controllers/v2/scrape.tsscrapeController