数据截至 (上游 commit 31522cf981bf)
爬取与异步任务: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 |
runNuqWorker | worker 主循环:取作业→执行→标记完成 | services/worker/nuq-worker-runner.ts |
crawl-logic.ts | finishCrawlSuper:收尾聚合 + 完成 webhook | services/worker/crawl-logic.ts |
| team-semaphore / concurrency-limit | 团队级并发限流 | services/worker/team-semaphore.ts、lib/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:40-52 定义了全部拒绝类型,含义如下:
| DenialReason | 什么时候拒绝 |
|---|---|
DEPTH_LIMIT | URL 路径段数 > 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_PROTOCOL | mailto:/tel:/ftp:/file: 等非 HTTP 协议 |
这份"带理由的拒绝"最终会出现在 crawl 的错误报告里,用户能知道为什么某个页没被爬。
3.3 过滤主入口:filterLinks
批量过滤走 filterLinks(crawler.ts:168,签名 filterLinks(sitemapLinks, limit, maxDepth, fromMap, skipRobots, ignoreDiscoveryDepth))。它的核心逻辑其实下沉到了 Rust:
真实实现在 crawler.ts:201,调用从 @mendable/firecrawl-rs 导入的 filterLinks(crawler.ts:28),把 excludes/includes/robots/allowBackwardCrawling 等一股脑传进去,让 Rust 侧做正则和 URL 判断,再把 Rust 返回的原始 denialReasons(如 "DEPTH_LIMIT")在 JS 侧翻译成上表那种带上下文的长句(crawler.ts:229-279,例如 DEPTH_LIMIT 会填入实际 depth 和 maxDepth)。
关键细节:有一条纯 JS 的回退路径。 如果 Rust 调用抛错,代码 catch 后落到 crawler.ts:304 起的手写 JS 过滤(逐条 new URL + 正则 + isRobotsAllowed + isFile)。这是 Firecrawl 反复出现的模式:Rust 快路径 + JS 慢回退,保证一条路挂了功能不崩。
filterURL(单条版,crawler.ts:715)同理,直接调 Rust 的 filterUrl,用于抓页后逐个判断新发现的链接。
3.4 文件类型判断:isFile
哪些扩展名算"文件、不爬",硬编码在 isFile(crawler.ts:837)的 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,分三步:
- 抓
getRobotsTxt(crawler.ts:468)通过fetchRobotsTxt拿${baseUrl}/robots.txt的内容; - 导入
importRobotsTxt(crawler.ts:518)用robots-parser解析,并读出 crawl-delay——它对多个 UA 名做兜底:getCrawlDelay(this.robotsUserAgent)失败就试"FireCrawlAgent"/"FirecrawlAgent"(crawler.ts:523-527);同时提取 robots 里声明的<sitemap>(crawler.ts:529); - 判定
isRobotsAllowed(crawler.ts:824)对单个 URL 用isUrlAllowedByRobots判断;若ignoreRobotsTxt为真则直接放行。
3.6 sitemap:比爬链接更高效的发现方式
顺着页面 <a> 一层层爬很慢;如果站点有 sitemap,一次就能拿到几千个 URL。tryGetSitemap(crawler.ts:544)就是干这个:
它会尝试多个 sitemap 位置(tryFetchSitemapLinks,crawler.ts:880),按优先级:
起始 URL 自身(若以 .xml 结尾则直接用,否则拼 /sitemap.xml)
├─ robots.txt 里声明的所有 sitemap
├─ 若是子域名 → 再试主域名的 /sitemap.xml(只保留回指子域的 URL)
└─ 兜底:baseUrl + /sitemap.xml
有个防爆炸的护栏:SITEMAP_LIMIT = 25(crawler.ts:31)——一次爬取最多命中 25 个 sitemap 文件,sitemapsHit 集合超过就停(crawler.ts:1074)。
sitemap 的解析在 WebScraper/sitemap.ts 的 getLinksFromSitemap(sitemap.ts:28):同样是 Rust 快路径 + JS 回退——先用 @mendable/firecrawl-rs 的 processSitemap(sitemap.ts:175),失败退到 parseSitemapXml,再失败退到 xml2js 的 parseStringPromise(sitemap.ts:198)。sitemap 索引(套娃 sitemap)会递归展开(instruction.action === "recurse",sitemap.ts:289)。注意 sitemap 内容本身是用抓取内核 scrapeURL 去下载的(sitemap.ts:105)——爬虫复用了单页抓取的引擎与容错。
3.7 抽链:从 HTML 到候选 URL
抓完一个页面后,要从它的 HTML 里挖出所有链接。extractLinksFromContent(crawler.ts:787)又是双路(markdown 内容会先走 extractLinksFromMarkdownContent 单独处理):
- Rust 路
extractLinksFromHTMLRust(crawler.ts:729):调extractLinks(Rust),再逐个过filterURL; - cheerio 回退
extractLinksFromHTMLCheerio(crawler.ts:741):Rust 挂了就用 cheerio 遍历$("a"),顺带还会解析<iframe>里data:text/html内联的 HTML(crawler.ts:760-771)。
4. 核心原理之二:爬取任务的生命周期
本节把 §2.3 那张图落到代码,追一次爬取从提交到收尾。
4.1 提交:crawlController
controllers/v2/crawl.ts 的 crawlController(crawl.ts:30)干的事,按顺序:
- 校验请求、检查权限、(可选)用 LLM 从自然语言 prompt 生成爬取参数(
crawl.ts:93-132); - 把
limit压到不超过剩余额度(crawl.ts:175); - 组装
StoredCrawl对象sc(crawl.ts:185),这是整个爬取的"配置 + 上下文"快照; crawlToCrawler(id, sc, flags)造一个WebCrawler(crawl.ts:209),并先抓一次 robots.txt 存进sc.robots(crawl.ts:212);- 解析队列后端、
crawlGroup.addGroup(...)建作业组(crawl.ts:223-233)——一次爬取的所有子作业都归在这个 group 下; saveCrawl(id, sc)存 Redis(crawl.ts:235)、markCrawlActive标记活跃(crawl.ts:237);- 入队一个
mode: "kickoff"作业(_addScrapeJobToBullMQ,crawl.ts:239),然后立刻给客户端返回 crawl id。
到此 HTTP 请求就返回了——真正的爬取在后台 worker 里进行。这就是"异步"的含义。
术语澄清:函数名叫
_addScrapeJobToBullMQ,但在此 commit 上底层队列已换成自研的 NuQ(见 §5);BullMQ 只是历史遗留的命名。
4.2 首轮扇出:processKickoffJob
kickoff 作业被 worker 取到后走 processKickoffJob(scrape-worker.ts:1228)。它是爬取的"点火器",负责铺开第一批要抓的页:
- 锁 + 入队起始 URL:
lockURL(scrape-worker.ts:1248)去重,再addScrapeJob入队起始页的single_urls作业,并标isCrawlSourceScrape: true(scrape-worker.ts:1266,标记它是爬取的"源头页"); - 探测 sitemap:若未禁用,拼出多个候选 sitemap 地址(起始页、
/sitemap.xml、根域 sitemap 等),对每个发一个kickoff_sitemap作业(addKickoffSitemapJob,scrape-worker.ts:1165/ 循环在1085); - 查索引补页:
kickoffGetIndexLinks(scrape-worker.ts:1135)从 Firecrawl 自己的 URL 索引里捞该站点已知的 URL,批量锁定并入队(scrape-worker.ts:1377-1434); - 最后
finishCrawlKickoff(scrape-worker.ts:1438)标记"点火阶段结束"——这是完成判定的一个必要条件(见 §4.5)。
kickoff_sitemap 作业单独由 processKickoffSitemapJob(scrape-worker.ts:1454)处理:用 scrapeSitemap 拉 sitemap → filterLinks 过滤 → 锁定 → 批量入队 single_urls,套娃 sitemap 再递归发 kickoff_sitemap。
4.3 递归扇出:processJob 抓完一页后做什么
普通页面作业走 processJob(scrape-worker.ts:355)。抓取本身委托给 startWebScraperPipeline(即 01 的内核,scrape-worker.ts:403)。抓完之后,才是爬取的精华——它就地发现并入队下一批页面(scrape-worker.ts:573-724):
# 示意,非源码:processJob 抓完一页后的扇出逻辑
crawler.setBaseUrl(该页最终 URL) # 处理重定向后的真实基准
links = crawler.filterLinks( # 抽链 + 过滤,一步到位
crawler.extractLinksFromContent(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:662-711),currentDiscoveryDepth 每扇出一层加 1(scrape-worker.ts:687)。抓完自己后,addCrawlJobDone(scrape-worker.ts:837)把本作业记入"完成集合"。
还有个重定向去重的巧处(scrape-worker.ts:548-570):如果页面 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>:visited | set | 已见过的 URL(去重) | lockURL (:449) |
crawl:<id>:visited_unique | set | 唯一已锁 URL,用于 limit 计数 | lockURL (:488) |
crawl:<id>:jobs | set | 本爬取的全部作业 id | addCrawlJob/getCrawlJobs (:114/:336) |
crawl:<id>:jobs_done | set | 已完成作业 id | addCrawlJobDone (:163) |
crawl:<id>:robots_blocked | set | 被 robots 拦的 URL | recordRobotsBlocked (:68) |
去重的核心是 lockURL(crawl-redis.ts:488): 它对 URL 先规范化(normalizeURL,crawl-redis.ts:383,可选去掉 query、统一 hash),再 SADD 到 visited;SADD 返回 1(新加入)才算"抢锁成功",这个作业才有资格入队。这是天然的分布式互斥——多个 worker 并发发现同一个 URL,只有一个能 SADD 到 1,其余得到 0、直接跳过。
lockURL 开头还查 limit(crawl-redis.ts:502-510):visited_unique 的基数 ≥ limit 就直接返回 false,让整张作业图停止生长。
URL 排列 generateURLPermutations(crawl-redis.ts:412) 是个精巧的去重放大器:同一个逻辑页可能有 http/https、www/非www、/、/index.html、/index.php 等写法。当开启 deduplicateSimilarURLs 时,lockURL 用规范排列的第一个变体做键(crawl-redis.ts:517),让这些等价 URL 折叠成一个。函数顶部那段注释(crawl-redis.ts:399-411)还写明了它必须满足的三条不变式,并说明由 permu-refactor.test.ts 证明。
4.5 收尾:什么时候算"爬完了"
判定条件在 isCrawlFinished(crawl-redis.ts:291),需要两件事同时成立:
爬取完成 ⟺ jobs_done 的数量 == jobs 的数量 (所有已知作业都跑完了)
AND kickoff 阶段已结束 + 所有 sitemap 作业已完成
第二个条件由 isCrawlKickoffFinished(crawl-redis.ts:303)保证——光是作业数相等还不够,必须确认"点火阶段没有还在扇出新作业",否则会在图还在生长时误判完成。
真正的收尾动作在 crawl-logic.ts 的 finishCrawlSuper(crawl-logic.ts:15),由 crawlFinishedQueue 触发(见 §5.5)。它 finishCrawl(crawl-redis.ts:338,标 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:终态,带returnvalue或failedReason;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)。crawlController 里 crawlGroup.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:330)后调 finishCrawlSuper(queue-worker.ts:233)完成 §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)
作业分发在 processJobInternal → processJobWithTracing(scrape-worker.ts:1628/1323)里按 mode 路由到 processKickoffJob / processKickoffSitemapJob / processJob。
团队并发限流分两层:
- 入队时的准入(
lib/concurrency-limit.ts+queue-jobs.ts):addScrapeJob(queue-jobs.ts:494)→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)。 - 完成时的放行(
concurrentJobDone,concurrency-limit.ts:307):一个作业跑完,腾出一个槽,就用getNextConcurrentJob(concurrency-limit.ts:215)从 backlog 里捞下一个提升为queued。捞取用ZPOPMIN原子弹出最小分数成员(concurrency-limit.ts:232),保证多个 worker 不会捞到同一个;捞出来若因"爬取级并发"仍不能跑,就先搁一边、最后再塞回(crawlBlocked,concurrency-limit.ts:259/276)。
同步 scrape 的并发走另一套:team-semaphore.ts 的 withSemaphore(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:47)不入队,而是:
- 用
teamConcurrencySemaphore.withSemaphore抢一个并发槽(scrape.ts:227); - 在API 进程里就地构造一个内存 NuQJob(
scrape.ts:258),关键是打了skipNuq: true(scrape.ts:288); - 直接
await processJobInternal(job)(scrape.ts:302)——同一个执行体,但没走 NuQ 的入队/领取/通知,而是就地同步跑完拿到 doc。
skipNuq: true 会让 processJobWithTracing 跳过并发槽的镜像维护、失败时直接抛错而非序列化回队列(scrape-worker.ts:1736/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:31)、每爬取最多 20 个 kickoff sitemap 作业(scrape-worker.ts:1172的TEMP: max 20,注释自认临时):超大站点可能漏发现部分 sitemap。 - 默认不向上爬:
allowBackwardCrawling关时,起始 URL 路径之外的页会被BACKWARD_CRAWLING拦掉——想爬整站得显式开crawlEntireDomain/allowBackwardCrawling。这常让用户困惑"为什么只爬到几页"。 - 去重是"尽力"而非绝对:
generateURLPermutations的注释坦承第 3 条不变式(不同 URL 的排列不重叠)"无法证明,超出爬虫范围"(crawl-redis.ts:408-410)——理论上存在极端 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 + 抽链;crawl是map的"发现"加上"抓取 + 递归扇出"。 - 横向(同 shelf,RAG/检索类):很多爬取框架用现成 Celery/BullMQ;Firecrawl 自研 NuQ 换取多租户公平并发 + 关系库持久化,是"规模优先"的取舍。
9. 代码地图(导航索引)
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 链接过滤主入口(Rust+JS 回退) | apps/api/src/scraper/WebScraper/crawler.ts | WebCrawler.filterLinks |
| 单条链接过滤 | apps/api/src/scraper/WebScraper/crawler.ts | WebCrawler.filterURL |
| 拒绝理由枚举 | apps/api/src/scraper/WebScraper/crawler.ts | DenialReason |
| 文件类型判断 | apps/api/src/scraper/WebScraper/crawler.ts | WebCrawler.isFile |
| robots.txt 抓取/导入/delay | apps/api/src/scraper/WebScraper/crawler.ts | getRobotsTxt / importRobotsTxt |
| sitemap 探测 | apps/api/src/scraper/WebScraper/crawler.ts | tryGetSitemap / tryFetchSitemapLinks |
| HTML 抽链(Rust/cheerio) | apps/api/src/scraper/WebScraper/crawler.ts | extractLinksFromContent |
| sitemap 解析 | apps/api/src/scraper/WebScraper/sitemap.ts | getLinksFromSitemap |
| 爬取提交控制器 | apps/api/src/controllers/v2/crawl.ts | crawlController |
| 首轮扇出 | apps/api/src/services/worker/scrape-worker.ts | processKickoffJob |
| sitemap 扇出 | apps/api/src/services/worker/scrape-worker.ts | processKickoffSitemapJob |
| 抓页后递归扇出 | apps/api/src/services/worker/scrape-worker.ts | processJob |
| 作业分发/skipNuq 分支 | apps/api/src/services/worker/scrape-worker.ts | processJobInternal / processJobWithTracing |
| 收尾聚合 + webhook | apps/api/src/services/worker/crawl-logic.ts | finishCrawlSuper |
| URL 去重锁 | apps/api/src/lib/crawl-redis.ts | lockURL |
| URL 等价排列 | apps/api/src/lib/crawl-redis.ts | generateURLPermutations |
| 完成判定 | apps/api/src/lib/crawl-redis.ts | isCrawlFinished / isCrawlKickoffFinished |
| NuQ 队列类 | apps/api/src/services/worker/nuq.ts | NuQ |
| 原子领取作业 | apps/api/src/services/worker/nuq.ts | getJobToProcess / prefetchJobs |
| 锁续期/完成/失败 | apps/api/src/services/worker/nuq.ts | renewLock / jobFinish / jobFail |
| 等待作业结果 | apps/api/src/services/worker/nuq.ts | waitForJob |
| 作业组 | apps/api/src/services/worker/nuq.ts | NuQJobGroup / crawlGroup |
| worker 主循环 | apps/api/src/services/worker/nuq-worker-runner.ts | runNuqWorker |
| 团队并发信号量(同步) | apps/api/src/services/worker/team-semaphore.ts | withSemaphore |
| 并发准入/放行(异步) | apps/api/src/lib/concurrency-limit.ts | getNextConcurrentJob / concurrentJobDone |
| 入队 + backlog | apps/api/src/services/queue-jobs.ts | addScrapeJob / addScrapeJobRaw |
| 同步 scrape 就地执行 | apps/api/src/controllers/v2/scrape.ts | scrapeController |