🎯 一款可定制、具备反检测功能的云浏览器,由自主研发的 Chromium驱动,专为网页爬虫AI 代理设计。👉立即试用
返回博客

在 Node.js 中构建分布式网络爬虫:队列和去重

Alex Johnson
Alex Johnson

Senior Web Scraping Engineer

29-Jul-2026

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 项目

该项目故意紧凑:

Copy
distributed-crawler/
├── package.json
└── src/
    └── crawler.mjs

依赖项版本在起草时与其包注册表进行了检查。该示例是一个先决条件缺口块:它需要 Node.js、Redis、一个 SCRAPELESS_API_KEY,以及一个在 ALLOWED_HOST 中设置的授权公共主机名。它设计为一个主机,因此队列限制器是完全特定于主机的。

json Copy
{
  "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"
  }
}

下面的工作者有四种终端结果:acceptedunchangedrejectedquarantined。它不会创建自动失败循环。操作员可以检查隔离记录,修复其原因,并显式地将规范 URL 重新入队。

javascript Copy
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。这些区分使仪表板更具可操作性。

按主机控制并发

工作程序并发回答:“此进程可以处理多少个作业?”主机请求预算回答:“此来源可以接收多少流量?”它们必须独立配置。

对于授权域:

  1. 读取 robots 规则和合同限制。
  2. 设置保守的每主机速度。
  3. 仅当队列和主机预算允许时运行多个 worker。
  4. 跟踪响应状态、内容接受和服务器延迟。
  5. 当来源显示压力或协议更改时,降低主机预算。

BullMQ workers 可以跨进程和机器共享队列。队列限制器协调所选队列,这就是示例使用特定于主机的队列的原因。对于许多域,从批准注册表生成队列并限制活跃工人对象的数量。

机器人排除协议由 RFC 9309 标准化。机器人规则不是许可的授权,并且不替代网站条款、隐私义务或适用法律。

委派页面获取而不失去控制

Node.js 控制层应决定可以收集什么。执行层应决定如何获取允许的页面表示。

Scrapeless Crawl 快速入门 文档记录异步爬取状态和页面级结果。对于需要更广泛交互或 JavaScript 执行的页面,请使用 Crawl 的浏览器选项,同时队列保持范围、作业身份和存储决策。

验证返回的内容,而不是假设执行者做出了业务决策。要求预期的文本、字段、语言、URL 和最小内容。在来源中存储获取路线,以便后期审计可以解释每个页面如何进入数据集。

在不超出范围的情况下发现链接

链接发现应在内容接受之后进行。仅解析满足内容合同的页面,然后在入队之前应用以下筛选:

  • 白名单主机名;
  • 允许的路径前缀;
  • 支持的 HTTP 方案;
  • 最大深度和页面计数;
  • 规范化和作业 ID 查找;
  • 文件类型排除;
  • 来源特定的查询参数政策。

不要让重定向悄无声息地扩展允许的主机集。记录最终 URL,将其与范围进行比较,并拒绝跨域结果,除非来源注册明确授权它们。

对于网站地图,将每个 URL 视为发现的输入,而不是可信的输出。通过与 HTML 链接相同的前沿路径进行规范化、过滤和去重。

存储页面版本和爬取来源

一个有用的存储模型有三个记录:

  1. 爬虫: 范围、种子 URL、截止日期、策略版本和整体状态。
  2. 页面: 规范键、源 URL、最新接受的哈希、采集时间和模式版本。
  3. 页面版本: 内容哈希、有效负载位置、元数据和获取来源。

队列不是长期数据库。作业清理策略可能会删除已完成的条目,而页面表必须根据保留政策保持幂等性键和版本历史。

如果页面未改变,更新新鲜度时间戳而不重复有效负载。如果有所更改,写入新的不可变页面版本,将页面记录指向它,并通过单独事件通知下游索引。

观察爬虫

仅凭队列长度可能会误导。爬虫可以在拒绝每个页面的同时清空其边界。追踪:

信号 回答的问题
每分钟接受的页面数 有用的数据是否到达?
发现到接受的延迟 流水线多久未更新?
重复入队率 规范化是否有效?
按原因拒绝率 源标记或验证是否发生变化?
隔离时间 操作债务是否在累积?
主机请求速度 政策是否被遵循?
内容变化率 刷新计划是否适当?
队列年龄百分位 工作者的能力是否足够?

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,我们仅访问公开可用的数据,并严格遵循适用的法律、法规和网站隐私政策。本博客中的内容仅供演示之用,不涉及任何非法或侵权活动。我们对使用本博客或第三方链接中的信息不做任何保证,并免除所有责任。在进行任何抓取活动之前,请咨询您的法律顾问,并审查目标网站的服务条款或获取必要的许可。

最受欢迎的文章

目录