如何在 BullMQ Worker 中识别验证码(Node.js)

BullMQ 验证码 worker 第一次跑起来通常是对的,然后会用两种很安静的方式让你失望。控制同时运行多少个任务的那个 worker 选项默认是 1,所以积压的识别任务只能一个一个地消化。而只要 processor 占住了事件循环,BullMQ 就会判定任务已经 stalled,把它交给另一个 worker,同一个验证码于是被识别了两次。这两件事都不是 bug。它们都是默认值,也都只要改一行。
你需要什么
- Redis,以及在运行 worker 的那个项目里装好 BullMQ。
- CapSkip 运行在一台 Windows 机器上,Node 客户端装在同一个项目里。
- sitekey 和页面 URL 通过任务数据传进来,而不是写死在代码里,这样一个队列就能服务所有表单。
- 在横向扩展之前先定好连接模式,因为 worker 最后往往会分布在比识别工具更多的机器上。
# npm install capskip npm install bullmq capskip
第 1 步:worker 跑在哪里,以及这需要哪种连接模式
这件事值得先定下来,因为 BullMQ 自己给的建议就是把一批 worker 分散到很多台不同的机器上运行,而你一旦这么做,回环地址的含义就和它在你笔记本上的含义不一样了。
连接模式有两种。本地模式绑定 127.0.0.1,只响应本机,当你的自动化程序和识别工具在同一台机器上时,这是正确的选择。服务器模式绑定你的内网地址或公网 IP,这样另一台机器、一台 VPS 或某个托管平台就能通过 API 访问到同一台 Windows 机器。服务器模式只改变识别工具监听哪个地址。它依然是你自己的硬件,也依然不按次计费。两种模式都位于 连接设置.
| worker 进程运行在哪里 | 用哪种连接模式 |
|---|---|
| 在 CapSkip 所在的机器上,单个进程 | Local 模式。127.0.0.1 在这里确实是对的 |
| 一批 worker 分布在你自己的内网里 | Server mode,填识别工具的内网地址 |
| worker 跑在容器里,或者跑在 VPS 上 | 服务器模式,配一个可达的地址和一条防火墙规则 |
对于跑在 VPS 上的 worker,值得去安排一个静态公网 IP,免得地址在部署运行的过程中变掉。把地址放进环境变量而不是源码里,因为你的笔记本和你的 worker 需要的值不一样。客户端不会自己读取任何环境变量,所以要在你的代码里读取 CAPSKIP_HOST 并传进去,就像下面的 worker 那样。
// npm install capskip
import { CapSkip } from "capskip";
// One client, shared by every job this worker handles.
// CAPSKIP_HOST is 127.0.0.1 locally and the solver's
// address on every other machine.
export const solver = new CapSkip({
host: process.env.CAPSKIP_HOST ?? "127.0.0.1",
port: 8080,
});第 2 步:concurrency 默认一次只处理一个任务
除非你明确设置,否则一个 BullMQ worker 一次只处理一个任务。这个默认值对 CPU 型工作是合理的,用在识别上却大错特错,因为一次识别几乎全是在等待。BullMQ 说得很直接:只有当 worker 执行异步操作时,比如调用数据库或外部 HTTP 服务,并发才有可能。通过 HTTP 轮询识别工具正是这种形态,所以事件循环全程都是空闲的。
这笔账值得算一次。如果一次 reCAPTCHA 识别要 20 秒,一个 worker 在默认值下每分钟消化 3 个任务。同一个 worker 把 concurrency 设成 20,在同样的硬件上每分钟消化 60 个,因为其中 19 个任务都卡在 socket 读取上,而不是在抢 CPU。
import { Worker } from "bullmq";
import { solver } from "./solver.js";
// A solve is an awaited HTTP call, so raising this costs
// almost no CPU. Size it for the solver machine and for
// what the target site will accept.
const worker = new Worker("captcha", async (job) => {
const { sitekey, pageUrl } = job.data;
const result = await solver.recaptcha(sitekey, pageUrl);
return await submitForm(pageUrl, result.code);
}, { connection: { host: "127.0.0.1", port: 6379 }, concurrency: 20 });在按次计费的识别服务上,这个数字其实是一个花钱的闸门,所以很多示例都把它设得很保守。在这里它是一个关于单台机器容量的问题,也是一个关于你要提交的那个站点愿意多快接受请求的问题。后者通常才是更紧的那个限制。
第 3 步:什么会让任务变成 stalled,以及为什么 stalled 意味着识别两次
这个故障会让人白白搭进去一个上午,因为日志看上去像是队列一切正常,而识别工具自己的日志却不是这么说的。
任务到达 worker 时,BullMQ 会给它加一把锁,别的东西就碰不到它了,同时 worker 必须不断告诉队列自己还在干活。这把锁就是 lockDuration,默认 30000 ms,worker 每隔这个间隔的一半续一次。另有一个单独的扫描 stalledInterval,每 30000 ms 运行一次,专门去找那些没人续过的锁。如果 worker 忙到来不及及时续锁,这个任务就会被标记为 stalled,退回 waiting,然后由另一个 worker 再处理一遍。一旦它 stalled 的次数超过 maxStalledCount 允许的上限,而这个值默认是 1,它就会转而进入 failed 集合。
所以 stalled 的任务不等于失败的任务,重新跑一遍也不是重试。这是队列在正确地假设:持有那个任务的 worker 已经死了。你的识别工具会看到同一个验证码被提交了两次,而真正被用上的只有第二个 token。
| Worker 选项 | 默认值 | 它决定什么 |
|---|---|---|
| "concurrency" | 1 | 一个 worker 同时处理多少个任务 |
| "lockDuration" | 30000 ms | 锁在没有续期的情况下能存活多久 |
| "lockRenewTime" | lockDuration 的一半 | worker 多久续一次锁 |
| "stalledInterval" | 30000 ms | 多久清理一次没有续期的锁 |
| "maxStalledCount" | 1 | 任务判定失败前允许重新跑多少次 |
原因永远是同一个:processor 占住了 CPU。一个被 await 的 HTTP 调用不会占住 CPU,所以客户端自己的轮询是安全的。手写一个在结果接口上忙等的循环就不安全了,在提交之前同步解码一张大图同样不安全。让客户端去轮询,因为它从 250 ms 起步并逐步退避,而不是按固定间隔一直睡,同时把任何真正 CPU 密集的步骤挪出 processor,或者挪进 sandboxed processor 里。
把锁的时长调大是错误的第一反应。那只是在治标,而且如果 worker 真的已经死了,一把很长的锁会让这个任务在整个窗口期里都没人碰。
第 4 步:重试,以及那个绝不能留着的 token
BullMQ 默认不重试。attempts 选项是 1,意思是试一次,然后进 failed 集合。对识别来说这通常是对的,因为这里大多数失败要么是参数写错了,那会一直以同样的方式失败,要么是站点已经改版了。重试真正有价值的场景,是 worker 短暂连不上识别工具,所以给它们配一个 backoff,并把次数保持得很小。
await queue.add("signup", { sitekey, pageUrl }, {
// Two tries, spaced out, for a solver that was
// briefly unreachable. Not for a bad sitekey.
attempts: 2,
backoff: { type: "exponential", delay: 5000 },
// Do not leave finished jobs sitting in Redis forever.
removeOnComplete: { age: 3600, count: 1000 },
});与此配套有两条规则。绝不要把一个 token 跨 attempt 带过去:一个 reCAPTCHA token 大约只有两分钟有效期,所以要在负责提交它的那一次 attempt 里面完成识别,另外也请读一读 这篇关于 token 过期的指南 ,如果你动了缓存 token 的念头的话。还有,不要把 token 作为任务的返回值。BullMQ 默认会保留已完成的任务,所以这个返回值会落进 Redis 并一直留在那里。应该返回提交的结果,那才是你之后真正想看的东西。
完整可运行示例
一个生产者、一个 worker,识别和消费它的那段代码放在同一个 processor 里。
// npm install capskip
import { Queue, Worker } from "bullmq";
import { CapSkip, ValidationException } from "capskip";
const connection = { host: "127.0.0.1", port: 6379 };
const queue = new Queue("captcha", { connection });
const solver = new CapSkip({
host: process.env.CAPSKIP_HOST ?? "127.0.0.1",
port: 8080,
});
new Worker("captcha", async (job) => {
const { sitekey, pageUrl, email } = job.data;
try {
// Solve and submit together. The token is short lived.
const result = await solver.recaptcha(sitekey, pageUrl);
const res = await postSignup(pageUrl, email, result.code);
return { status: res.status };
} catch (err) {
// A bad sitekey fails the same way on every attempt.
if (err instanceof ValidationException) {
await job.discard();
}
throw err;
}
}, { connection, concurrency: 20 });这个调用针对的是 reCAPTCHA v2。其他类型的写法完全一样:把 invisible 或 enterprise 设为 1,或者把 version 设为 v3 并带上一个 action,又或者改为调用 turnstile 或 geetest。完整的接口请见 Node.js 验证码识别页面.
常见错误及其含义
| 你所看到的 | 原因 | 修复 |
|---|---|---|
| 队列消化的速度远远慢于识别工具能达到的速度 | worker 的 concurrency 还停在默认值 1 | 调高它。识别是被 await 的 I/O,不是 CPU 运算 |
| 同一个任务记录了两次识别,相隔几秒 | 锁没有续上,所以任务被算作 stalled | 不要在 processor 里阻塞事件循环 |
| 任务进了 failed 集合,却没有抛出任何错误 | 它 stalled 的次数超过了允许的上限 | 同样是那段阻塞的代码。要修的是它,不是计数器 |
| 在你笔记本上正常,在 worker 机器上报 NetworkException | 那台机器上的 127.0.0.1 指的就是那台机器 | 用服务器模式,并为 worker 设置 CAPSKIP_HOST |
| 只有在第二次 attempt 时 token 被拒绝 | 第一次 attempt 的 token 被带到了后面 | 在负责提交的那一次 attempt 里完成识别 |
| ApiException 里带着 ERROR_WRONG_USER_KEY | worker 运行的那个环境里没有设置 CAPSKIP_API_KEY | 在 worker 运行的地方设置这个变量,然后重启那个 worker |
| 手写轮询时返回 CAPCHA_NOT_READY | 结果还没完成就被读取了 | 交给客户端轮询,它会自己退避 |
最后那个响应就是这么拼的,少掉的那个字母不是我们这边的笔误,因为 API 返回的确实就是这个写法。完整解释见 一篇关于 CAPCHA_NOT_READY 响应的完整说明.
常见问题
我的 worker 可以和识别工具跑在不同的机器上吗?
可以,而且一旦你扩展到不止一个进程,这就是常规做法。在连接设置里把 CapSkip 切换到服务器模式,让它监听你的网络地址而不是回环地址,然后在每个 worker 的环境里设置 CAPSKIP_HOST。在你自己的网络里,这个地址就是一个内网地址,不需要任何东西暴露到公网。如果某个 worker 在 VPS 上,就用一个静态公网 IP,再加一条只允许你预期的那些地址的防火墙规则。
concurrency 到底该设成多少?
从 10 开始,然后盯住两件事:识别工具所在的机器,以及你要提交的那个站点的反应。因为一次识别是被 await 的 I/O,worker 进程本身很少成为瓶颈。通常是站点先撑不住,它会用限流来告诉你,而这远在 Node 耗尽余地之前就会发生。这里没有任何东西要排队等余额,所以这个数字是一个容量决策,而不是预算决策。
我从来没设置过 attempts,为什么同一个验证码被识别了两次?
因为那次重新运行不是重试。stalled 的任务不管 attempts 怎么设都会被重新排队,前提假设是持有它的 worker 已经死了。触发条件是锁没有及时续期,而这发生在 processor 占住 CPU、而不是在 await 的时候。找到 processor 里那段同步代码并把它挪走,重复识别就消失了。
识别应该独立成一个队列,让其他任务来调用吗?
通常不该。拆开意味着 token 要跨过一个队列边界,并在父任务恢复期间躺在 Redis 里,而这是最快让你用上一个已经过期的 token 的办法。把识别和消费它的代码放在同一个 processor 里,返回结果而不是返回凭据。只有当一个独立队列交回来的根本不是 token 时,它才说得通。
简短版结论
把 worker 的 concurrency 调高,因为一次只跑一个任务的默认值毫无道理地卡住了一整队被 await 的 HTTP 调用。让 processor 不要占住 CPU,锁才能持续续期,因为 stalled 的任务会被另一个 worker 重新跑一遍,而你要用白白浪费掉的吞吐来为此买单。把 attempts 留在很低的值,并且绝不要把 token 跨 attempt 带过去。只要有 worker 跑在别的地方,就把地址放进 CAPSKIP_HOST,并把识别工具切换到服务器模式。
- 客户端背后的原始端点,文档见 CapSkip API 文档.
- 勾选框挑战本身的说明,请见 reCAPTCHA v2 识别页面.
在你确定 concurrency 这个数字之前,有一点值得掂量:CapSkip 是一个 无限量验证码识别工具 ,它运行在你已经拥有的硬件上,所以把这个数字调高,消耗的是你一台机器的容量,而不会产生任何按次费用。
