在 Node.js 中构建分布式网络爬虫:队列和去重
Senior Web Scraping Engineer
TL;DR:
- 一个分布式网络爬虫需要一个持久化的 URL 边界、无状态的工作节点、规范的 URL 键、每个主机的请求预算以及明确的终端隔离路径。
- 让 Node.js 负责发现、调度、状态和存储。将 JavaScript 渲染和页面获取委托给一个受管理的执行层,必要时对源进行处理。
- 使用规范的 URL 哈希作为队列作业 ID 和存储幂等密钥。这可以在重复工作到达工作节点之前阻止它。
- 全局工作节点并发和每个主机的节奏解决不同的问题。根据容量扩展第一个;根据源的许可和观察到的服务器行为设置第二个。
- 从一个小的授权源集开始,测量接受的页面而不是尝试的页面,并在队列深度、新鲜度延迟和架构拒绝可见后进行扩展。
单进程爬虫以可预测的方式失败:其内存队列在重启时消失,重复链接成倍增加,一个缓慢的主机占用了事件循环,而浏览器执行消耗了应该调度工作的同一台机器。
分布式网络爬虫将这些职责分开。Node.js 拥有控制平面——URL 边界、作业状态、去重和存储决策。独立的工作节点拥有数据路径。对于 JavaScript 渲染的公共页面,一个管理服务可以执行页面并返回 Markdown 或 HTML,而不会将浏览器进程放置在每个工作节点容器内。
定义要求和失败模式
在选择队列之前,编写爬虫合同:
- 哪些域和路径是授权的?
- 每个爬虫可以发现多少页面和层级?
- 每个源需要多大的新鲜度窗口?
- 哪些响应格式和必填字段定义被接受的页面?
- 每个主机允许什么请求节奏?
- 无效、空或意外的结果会去哪里?
- 作业和页面版本必须保留多长时间?
第一个版本还应具有停止条件:最大深度、最大接受页面、最大发现的 URL 和截止日期。这些限制防止日历归档、分面导航或跟踪参数将小型爬虫变成一个开放式图形遍历。
常见的失败模式是架构性的:
| 失败 | 根本原因 | 控制 |
|---|---|---|
| 边界丢失 | 内存队列 | 基于 Redis 的持久作业 |
| 重复页面 | 使用原始 URL 作为键 | 规范的 URL 哈希 |
| 一台主机过载 | 仅设置了全局并发 | 每个主机队列和节奏 |
| 存储空内容 | 传播成功被视为数据成功 | 内容接收合同 |
| 工作节点无法扩展 | 本地浏览器状态 | 无状态工作节点和受管理的执行 |
| 毒药作业无限循环 | 没有终端状态 | 带手动释放的隔离队列 |
| 爬虫看起来正常但陈旧 | 仅测量吞吐量 | 新鲜度延迟和接受页面指标 |
将控制平面与页面执行分开
爬虫可以被描绘为两个平面:
控制平面: 种子 → URL 规范化 → 去重 → Redis 队列 → 作业状态 → 存储元数据
执行平面: 工作节点 → 页面获取 → 内容验证 → 链接提取 → 接受记录或隔离
这个边界很重要,因为调度和浏览器执行的扩展方式不同。队列操作是小而有状态的。页面执行网络密集,可能需要 JavaScript、区域路由或孤立的浏览器会话。
Scrapeless Crawl 支持单页面、批处理和链接网站的收集,其格式包括 Markdown、HTML、链接、元数据和截图。当执行路径需要交互式浏览器时,Scrapeless Scraping Browser 将该运行时保持在工作节点容器外。Node 应用程序可以保持为范围和状态的记录系统,而 Scrapeless 则处理获取。
要了解发现和提取的概念介绍,请阅读 什么是网络爬虫。以下设计从单机爬虫停止的地方开始:共享队列、幂等入队和多个工作节点。
构建基于 Redis 的 URL 边界
BullMQ 在 Redis 上实现分布式作业执行。它的 官方文档 涵盖了队列、工作节点、事件、延迟作业、速率限制和作业状态。
边界应该存储一个小的作业负载:
| 字段 | 目的 |
|---|---|
url |
要收集的规范 URL |
host |
队列分区和策略查找 |
depth |
发现边界 |
crawlId |
相关性与取消 |
parentUrl |
发现的来源 |
schemaVersion |
下游合同 |
不要将页面 HTML 存储在 Redis 中。将大型结果存储在对象存储或数据库中,仅在队列中保留标识符、状态和简洁的元数据。
当不同的域需要不同的请求预算时,每个批准的主机使用单独的队列。分配给 docs-example-com 的工作者可以使用一个限制器,而另一个主机则获得自己的节奏。全局队列更简单,但其限制器无法表达独立的主机策略。
使入队操作幂等
两个页面可以是等效的,但其原始 URL 不同:
- 一个片段指向同一文档内部的一个位置;
- 跟踪参数改变而内容不变;
- 查询参数的顺序不同;
- 默认端口或尾随斜杠各异;
- 相对链接解析为同一绝对页面。
在队列插入之前进行归一化。WHATWG URL 标准 定义了 Node.js URL 类实现的解析模型。
安全策略是特定于源的。删除每一个查询参数可能会合并真正不同的页面。维护一个已知参数的允许列表或拒绝列表,并保留任何改变资源的参数。
在归一化之后,使用 SHA-256 对 URL 进行哈希处理。将该十六进制摘要用作:
- BullMQ 作业 ID;
- 数据库的 upsert 键;
- 作业与存储页面之间的来源链接。
当队列中已存在具有相同 ID 的作业时,BullMQ 会忽略新作业。保留策略删除的作业记录不再提供该保护,因此持久存储必须保持相同的幂等性密钥。
最小版本固定的 Node.js 项目
该项目故意紧凑:
distributed-crawler/
├── package.json
└── src/
└── crawler.mjs
依赖项版本在起草时与其包注册表进行了检查。该示例是一个先决条件缺口块:它需要 Node.js、Redis、一个 SCRAPELESS_API_KEY,以及一个在 ALLOWED_HOST 中设置的授权公共主机名。它设计为一个主机,因此队列限制器是完全特定于主机的。
json
{
"name": "distributed-crawler-example",
"private": true,
"type": "module",
"scripts": {
"start": "node src/crawler.mjs"
},
"dependencies": {
"@scrapeless-ai/sdk": "1.3.1",
"bullmq": "5.79.3"
},
"engines": {
"node": ">=22"
}
}
下面的工作者有四种终端结果:accepted、unchanged、rejected 和 quarantined。它不会创建自动失败循环。操作员可以检查隔离记录,修复其原因,并显式地将规范 URL 重新入队。
javascript
import { createHash } from "node:crypto";
import { Queue, Worker } from "bullmq";
import { ScrapingCrawl } from "@scrapeless-ai/sdk";
const redis = {
host: process.env.REDIS_HOST ?? "127.0.0.1",
port: Number(process.env.REDIS_PORT ?? 6379)
};
const allowedHost = process.env.ALLOWED_HOST ?? "example.com";
const queueName = `crawl-${hostKey(allowedHost)}`;
const frontier = new Queue(queueName, { connection: redis });
const quarantine = new Queue(`${queueName}-quarantine`, { connection: redis });
const crawl = new ScrapingCrawl({
apiKey: process.env.SCRAPELESS_API_KEY
});
function sha256(value) {
return createHash("sha256").update(value).digest("hex");
}
function hostKey(host) {
return host.toLowerCase().replaceAll(".", "-");
}
function canonicalize(input) {
const url = new URL(input);
url.hash = "";
url.hostname = url.hostname.toLowerCase();
for (const key of ["utm_source", "utm_medium", "utm_campaign"]) {
url.searchParams.delete(key);
}
url.searchParams.sort();
return url.href;
}
async function enqueue(url, crawlId, depth = 0, parentUrl = null) {
const canonicalUrl = canonicalize(url);
const parsed = new URL(canonicalUrl);
if (parsed.hostname !== allowedHost) {
throw new Error(`Host outside crawl scope: ${parsed.hostname}`);
}
const key = sha256(canonicalUrl);
await frontier.add(
"collect-page",
{
url: canonicalUrl,
host: parsed.hostname,
depth,
crawlId,
parentUrl,
schemaVersion: "crawl-page-v1"
},
{
jobId: key,
removeOnComplete: 1000,
removeOnFail: false
}
);
return key;
}
const worker = new Worker(
queueName,
async (job) => {
try {
const result = await crawl.scrapeUrl(job.data.url, {
formats: ["markdown", "links"],
onlyMainContent: true,
timeout: 15000
});
const markdown = result.markdown ?? result.data?.markdown;
if (typeof markdown !== "string" || markdown.length < 200) {
return { state: "rejected", reason: "content-contract", url: job.data.url };
}
const record = {
key: job.id,
url: job.data.url,
crawlId: job.data.crawlId,
depth: job.data.depth,
contentHash: sha256(markdown),
markdown
};
console.log(JSON.stringify({ state: "accepted", ...record }));
return { state: "accepted", key: record.key, contentHash: record.contentHash };
} catch (error) {
await quarantine.add("inspect-page", {
...job.data,
sourceJobId: job.id,
reason: error instanceof Error ? error.message : "unknown"
});
return { state: "quarantined", key: job.id };
}
},
{
connection: redis,
concurrency: 4,
limiter: { max: 2, duration: 1000 }
}
);
worker.on("completed", (job, result) => {
console.log(JSON.stringify({ event: "completed", jobId: job.id, result }));
});
worker.on("failed", (job, error) => {
console.error(JSON.stringify({
event: "worker-failed",
jobId: job?.id,
message: error.message
}));
});
await enqueue(`https://${allowedHost}/`, "demo-crawl");
这个示例打印接受的记录,以使数据合同可见。将 console.log 替换为以 record.key 作为键的数据库插入。在写入新页面版本或重建索引之前,比对 contentHash 和存储的值。
理解任务状态机
爬虫需要一个状态模型,操作人员能够解释:
| 状态 | 意义 | 下一步操作 |
|---|---|---|
queued |
规范 URL 正在等待 | Worker 声明它 |
active |
一个 worker 拥有租约 | 获取并验证 |
accepted |
内容通过了合同 | 存储、索引、发现链接 |
unchanged |
内容哈希与存储版本匹配 | 更新新鲜度元数据 |
rejected |
响应完成但内容无效 | 审查验证器或来源 |
quarantined |
执行未能产生决定 | 手动检查和释放 |
cancelled |
爬取范围或截止日期结束 | 保留审计元数据 |
将 failed 作为队列报告的基础设施事件,而不是唯一业务状态。返回空应用程序外壳的作业在技术上是完成的,但应被标记为 rejected。在获取之前,应将超出范围的页面标记为 cancelled。这些区分使仪表板更具可操作性。
按主机控制并发
工作程序并发回答:“此进程可以处理多少个作业?”主机请求预算回答:“此来源可以接收多少流量?”它们必须独立配置。
对于授权域:
- 读取 robots 规则和合同限制。
- 设置保守的每主机速度。
- 仅当队列和主机预算允许时运行多个 worker。
- 跟踪响应状态、内容接受和服务器延迟。
- 当来源显示压力或协议更改时,降低主机预算。
BullMQ workers 可以跨进程和机器共享队列。队列限制器协调所选队列,这就是示例使用特定于主机的队列的原因。对于许多域,从批准注册表生成队列并限制活跃工人对象的数量。
机器人排除协议由 RFC 9309 标准化。机器人规则不是许可的授权,并且不替代网站条款、隐私义务或适用法律。
委派页面获取而不失去控制
Node.js 控制层应决定可以收集什么。执行层应决定如何获取允许的页面表示。
Scrapeless Crawl 快速入门 文档记录异步爬取状态和页面级结果。对于需要更广泛交互或 JavaScript 执行的页面,请使用 Crawl 的浏览器选项,同时队列保持范围、作业身份和存储决策。
验证返回的内容,而不是假设执行者做出了业务决策。要求预期的文本、字段、语言、URL 和最小内容。在来源中存储获取路线,以便后期审计可以解释每个页面如何进入数据集。
在不超出范围的情况下发现链接
链接发现应在内容接受之后进行。仅解析满足内容合同的页面,然后在入队之前应用以下筛选:
- 白名单主机名;
- 允许的路径前缀;
- 支持的 HTTP 方案;
- 最大深度和页面计数;
- 规范化和作业 ID 查找;
- 文件类型排除;
- 来源特定的查询参数政策。
不要让重定向悄无声息地扩展允许的主机集。记录最终 URL,将其与范围进行比较,并拒绝跨域结果,除非来源注册明确授权它们。
对于网站地图,将每个 URL 视为发现的输入,而不是可信的输出。通过与 HTML 链接相同的前沿路径进行规范化、过滤和去重。
存储页面版本和爬取来源
一个有用的存储模型有三个记录:
- 爬虫: 范围、种子 URL、截止日期、策略版本和整体状态。
- 页面: 规范键、源 URL、最新接受的哈希、采集时间和模式版本。
- 页面版本: 内容哈希、有效负载位置、元数据和获取来源。
队列不是长期数据库。作业清理策略可能会删除已完成的条目,而页面表必须根据保留政策保持幂等性键和版本历史。
如果页面未改变,更新新鲜度时间戳而不重复有效负载。如果有所更改,写入新的不可变页面版本,将页面记录指向它,并通过单独事件通知下游索引。
观察爬虫
仅凭队列长度可能会误导。爬虫可以在拒绝每个页面的同时清空其边界。追踪:
| 信号 | 回答的问题 |
|---|---|
| 每分钟接受的页面数 | 有用的数据是否到达? |
| 发现到接受的延迟 | 流水线多久未更新? |
| 重复入队率 | 规范化是否有效? |
| 按原因拒绝率 | 源标记或验证是否发生变化? |
| 隔离时间 | 操作债务是否在累积? |
| 主机请求速度 | 政策是否被遵循? |
| 内容变化率 | 刷新计划是否适当? |
| 队列年龄百分位 | 工作者的能力是否足够? |
OpenTelemetry 在其 遥测信号指南 中描述了跟踪、指标、日志和行李。使用一个爬虫 ID 在生成者、队列事件、获取调用、存储写入和下游索引事件中。
对陈旧的接受数据和旧的隔离条目发出警报,而不仅仅是在进程崩溃时。一个进程可以是健康的,而数据集却悄然停止变化。
部署清单
- 运行具有持久性、身份验证、网络控制和适应工作负载的备份的 Redis。
- 保持工作者无状态,并在各实例之间部署相同的镜像。
- 将 API 密钥存储在秘密管理器或环境注入中,绝不要在作业有效负载中。
- 单独定义队列保留与页面保留。
- 限制爬虫深度、页面数、时间和主机范围。
- 使用生产者和工作者共享的一个主机政策来源。
- 在部署期间排空工作者,以免活跃租约被放弃。
- 测试取消、隔离释放、存储幂等性和源政策变更。
- 记录 Node.js、BullMQ、Scrapeless SDK 和页面模式的版本。
从一个获得批准的主机和一个小的页面上限开始。只有在仪表板显示接受的数据、新鲜度和请求政策保持在目标范围内后,才添加域名。
结论:扩展边界,而不是不确定性
当每个 URL 具有一个规范身份、每个主机都有明确的预算,且每个页面以有意义的状态结束时,分布式网络爬虫变得可靠。Node.js 和 BullMQ 可以负责边界和工作者协调;Scrapeless 可以负责需要管理渲染和网络处理的页面执行。
创建一个有界的概念验证,使用一个公共的、授权的域名,查看 Scrapeless 定价页面,然后 创建一个 Scrapeless 账户。通过 Crawl 路由获取,并在增加工作者数量之前测量接受的页面和新鲜度。
常见问题解答
什么使网络爬虫分布式?
它的 URL 边界和作业状态在多个工作者进程或机器之间共享。工作者可以声claim独立作业,将结果写入共享存储,并且可以在不依赖于一个进程的内存的情况下扩展。
为什么对 Node.js 爬虫使用 BullMQ?
BullMQ 提供基于 Redis 的队列、分布式工作者、并发控制、事件和作业标识符。爬虫仍然需要自己的 URL 策略、内容契约、持久存储和可观察性。
工作者之间的 URL 去重是如何工作的?
在入队之前标准化 URL,哈希规范形式,并使用摘要作为队列作业 ID 和存储键。Redis 协调作业创建,而数据库在队列保留移除旧作业后保持幂等性。
每个工作者应该启动自己的浏览器吗?
不一定。本地浏览器会增加容器大小、内存使用和操作工作。一个管理的获取层可以返回渲染的内容,而工作者则专注于调度、验证、发现和存储。
应该如何处理失败的页面作业?
将业务结果与基础设施事件分开。可以拒绝无效内容;未解决的执行错误可以进入一个隔离队列进行检查和显式释放。避免无限自动循环。
在Scrapeless,我们仅访问公开可用的数据,并严格遵循适用的法律、法规和网站隐私政策。本博客中的内容仅供演示之用,不涉及任何非法或侵权活动。我们对使用本博客或第三方链接中的信息不做任何保证,并免除所有责任。在进行任何抓取活动之前,请咨询您的法律顾问,并审查目标网站的服务条款或获取必要的许可。



