DevOps & Scaling

验证码解决工作队列的滚动更新

给验证码识别 worker 集群发版,最稳的做法是:一次只动一个节点,先排空在途任务,再换版本,健康检查通过才轮到下一个。全程都有容量在跑,调用方无感。

直接全部重启的代价很实在:验证码识别是长耗时异步任务,提交到 in.phpcaptcha_id 再轮询 res.php,一次 reCAPTCHA v2 常要几十秒。进程被杀掉时,这些已提交、已计费的任务就成了孤儿。

发版前先定三个参数

maxUnavailable(同时离线几个节点):10 个以内填 1,再大按总数 10%–25% 取值,不要超过 50%。

drain_timeout(排空等多久):取“最慢那类验证码的识别耗时 + 缓冲”。reCAPTCHA v2 一般 30–90 秒,设 120 秒即可。

健康检查判定方式:“进程还活着”最省事也最没用。让新节点走一次完整 API 往返——API Key 没注入、依赖装错版本,只有真实往返才暴露。

节点状态流转

Workers: [W1-old] [W2-old] [W3-old] [W4-old]

Step 1:  [W1-drain] [W2-old]  [W3-old]  [W4-old]
Step 2:  [W1-NEW✓]  [W2-old]  [W3-old]  [W4-old]
Step 3:  [W1-NEW✓]  [W2-drain] [W3-old]  [W4-old]
Step 4:  [W1-NEW✓]  [W2-NEW✓]  [W3-old]  [W4-old]
  ...until all updated

关键是中间的 drain 状态:节点排空后不再接新任务,手上的继续跑完。路由层只把任务发给 RUNNING 节点。

Python:带健康门禁与自动回滚的协调器

Worker 管单节点的提交、轮询与排空,RollingUpdateOrchestrator 按批推进并在检查不过时回滚。

import os
import time
import signal
import threading
import requests
from dataclasses import dataclass, field
from enum import Enum

API_KEY = os.environ["CAPTCHAAI_API_KEY"]


class WorkerState(Enum):
    RUNNING = "running"
    DRAINING = "draining"
    STOPPED = "stopped"
    UPDATING = "updating"


@dataclass
class Worker:
    worker_id: str
    version: str
    state: WorkerState = WorkerState.RUNNING
    active_tasks: int = 0
    tasks_completed: int = 0
    session: requests.Session = field(default_factory=requests.Session)

    def solve(self, task):
        if self.state != WorkerState.RUNNING:
            return {"error": "WORKER_NOT_ACCEPTING"}

        self.active_tasks += 1
        try:
            result = self._do_solve(task)
            self.tasks_completed += 1
            return result
        finally:
            self.active_tasks -= 1

    def _do_solve(self, task):
        resp = self.session.post("https://ocr.captchaai.com/in.php", data={
            "key": API_KEY,
            "method": task.get("method", "userrecaptcha"),
            "googlekey": task["sitekey"],
            "pageurl": task["pageurl"],
            "json": 1
        })
        data = resp.json()
        if data.get("status") != 1:
            return {"error": data.get("request")}

        captcha_id = data["request"]
        for _ in range(60):
            time.sleep(5)
            result = self.session.get(
                "https://ocr.captchaai.com/res.php",
                params={
                    "key": API_KEY,
                    "action": "get",
                    "id": captcha_id,
                    "json": 1
                }
            ).json()
            if result.get("status") == 1:
                return {"solution": result["request"]}
            if result.get("request") != "CAPCHA_NOT_READY":
                return {"error": result.get("request")}
        return {"error": "TIMEOUT"}

    def drain(self, timeout=120):
        """Stop accepting tasks and wait for active tasks to complete."""
        self.state = WorkerState.DRAINING
        start = time.time()
        while self.active_tasks > 0:
            if time.time() - start > timeout:
                print(f"Worker {self.worker_id}: drain timeout with "
                      f"{self.active_tasks} tasks remaining")
                break
            time.sleep(1)
        self.state = WorkerState.STOPPED

    @property
    def is_healthy(self):
        return self.state == WorkerState.RUNNING


class RollingUpdateOrchestrator:
    def __init__(self, workers):
        self.workers = {w.worker_id: w for w in workers}
        self.lock = threading.Lock()

    def get_available_worker(self):
        """Route tasks only to RUNNING workers."""
        with self.lock:
            for worker in self.workers.values():
                if worker.state == WorkerState.RUNNING:
                    return worker
        return None

    def rolling_update(self, new_version, health_check_fn=None,
                       max_unavailable=1, drain_timeout=120):
        """Update workers one at a time with health gates."""
        worker_ids = list(self.workers.keys())
        updated = []
        failed = []

        for i in range(0, len(worker_ids), max_unavailable):
            batch = worker_ids[i:i + max_unavailable]

            for wid in batch:
                worker = self.workers[wid]
                print(f"[{wid}] Draining (v{worker.version})...")

                # Step 1: Drain active tasks
                worker.drain(timeout=drain_timeout)

                # Step 2: "Deploy" new version
                print(f"[{wid}] Deploying v{new_version}...")
                worker.state = WorkerState.UPDATING
                worker.version = new_version
                time.sleep(2)  # Simulate deployment

                # Step 3: Start and health check
                worker.state = WorkerState.RUNNING
                if health_check_fn:
                    healthy = health_check_fn(worker)
                    if not healthy:
                        print(f"[{wid}] Health check FAILED — rolling back")
                        failed.append(wid)
                        self._rollback(updated)
                        return {
                            "status": "rolled_back",
                            "failed_at": wid,
                            "updated": updated,
                        }

                updated.append(wid)
                print(f"[{wid}] Updated to v{new_version} ✓")

        return {"status": "complete", "updated": updated, "failed": failed}

    def _rollback(self, updated_ids):
        """Roll back already-updated workers."""
        for wid in updated_ids:
            worker = self.workers[wid]
            print(f"[{wid}] Rolling back...")
            worker.state = WorkerState.STOPPED
            time.sleep(1)
            worker.version = "rollback"
            worker.state = WorkerState.RUNNING

    @property
    def status(self):
        return {
            wid: {
                "version": w.version,
                "state": w.state.value,
                "active_tasks": w.active_tasks,
            }
            for wid, w in self.workers.items()
        }


# Create fleet
workers = [Worker(f"w{i}", "1.2.0") for i in range(6)]
orchestrator = RollingUpdateOrchestrator(workers)


def health_check(worker):
    """Verify worker can solve a test CAPTCHA."""
    # In production, send a real test task
    return worker.state == WorkerState.RUNNING


# Execute rolling update
result = orchestrator.rolling_update(
    new_version="1.3.0",
    health_check_fn=health_check,
    max_unavailable=1,
    drain_timeout=60
)
print(f"Rolling update result: {result}")

几处细节:get_available_worker() 加锁避免路由与状态切换打架;drain() 超时只告警不抛异常;检查失败直接 _rollback()——失败要停在第一个节点上

Node.js:带进度统计与失败阈值

逻辑一致,多了进度统计和“失败率超 25% 中止”的保险丝。

const axios = require("axios");

const API_KEY = process.env.CAPTCHAAI_API_KEY;

class RollingUpdater {
  constructor(workerCount, currentVersion) {
    this.workers = Array.from({ length: workerCount }, (_, i) => ({
      id: `worker-${i}`,
      version: currentVersion,
      state: "running",
      activeTasks: 0,
    }));
    this.progress = { total: workerCount, completed: 0, failed: 0 };
  }

  async update(newVersion, options = {}) {
    const {
      maxUnavailable = 1,
      drainTimeout = 60000,
      healthCheckRetries = 3,
    } = options;

    console.log(
      `Starting rolling update: v${this.workers[0].version} → v${newVersion}`
    );

    for (let i = 0; i < this.workers.length; i += maxUnavailable) {
      const batch = this.workers.slice(i, i + maxUnavailable);

      for (const worker of batch) {
        try {
          // Drain
          console.log(`[${worker.id}] Draining...`);
          worker.state = "draining";
          await this.waitForDrain(worker, drainTimeout);

          // Deploy
          console.log(`[${worker.id}] Deploying v${newVersion}...`);
          worker.state = "updating";
          worker.version = newVersion;

          // Health check
          worker.state = "running";
          const healthy = await this.healthCheck(worker, healthCheckRetries);

          if (!healthy) {
            worker.state = "failed";
            this.progress.failed++;
            console.log(`[${worker.id}] FAILED health check`);

            if (this.progress.failed > Math.floor(this.workers.length * 0.25)) {
              console.log("Too many failures — aborting rolling update");
              return { status: "aborted", progress: this.progress };
            }
            continue;
          }

          this.progress.completed++;
          console.log(
            `[${worker.id}] Updated ✓ (${this.progress.completed}/${this.progress.total})`
          );
        } catch (err) {
          console.error(`[${worker.id}] Error: ${err.message}`);
          this.progress.failed++;
        }
      }
    }

    return { status: "complete", progress: this.progress };
  }

  async waitForDrain(worker, timeout) {
    const start = Date.now();
    while (worker.activeTasks > 0 && Date.now() - start < timeout) {
      await new Promise((r) => setTimeout(r, 1000));
    }
  }

  async healthCheck(worker, retries) {
    for (let attempt = 0; attempt < retries; attempt++) {
      try {
        const resp = await axios.get("https://ocr.captchaai.com/res.php", {
          params: { key: API_KEY, action: "getbalance", json: 1 },
          timeout: 10000,
        });
        if (resp.data.status === 1) return true;
      } catch {
        // Retry
      }
      await new Promise((r) => setTimeout(r, 5000));
    }
    return false;
  }
}

// Execute
const updater = new RollingUpdater(8, "1.2.0");
updater
  .update("1.3.0", { maxUnavailable: 2, drainTimeout: 30000 })
  .then((result) => console.log("Result:", JSON.stringify(result, null, 2)));

healthCheck 的重试很有必要:新节点刚起来时连接池和 DNS 没热,首次失败不代表版本有问题。

排障对照表

现象 原因 处理方式
更新期间任务被丢弃 drain_timeout 太短 调大超时,核对最长任务耗时
健康检查一直不过 新版本本身有问题 先回滚,放到 staging 验证再发
整轮更新拖太久 maxUnavailable 太小 提到 2–3,或按 10%–25% 取值
回滚后版本不统一 部分回滚没做完 按已更新列表全量回退

几种发布策略怎么选

发布策略 停机 回滚速度 实现复杂度 适用场景
滚动更新 中等 多数发版
蓝绿部署 立即 核心链路
金丝雀发布 50 个节点以上
停机重建 有,短暂 不适用 最低 开发测试环境

能接受“逐个回退花几分钟”就用滚动更新;要秒级回退就上蓝绿部署

一个国内团队常见的场景

跨境采集团队的集群常混着两类任务:国际站点的 reCAPTCHA v2 与 Cloudflare Turnstile,国内站点的 GeeTest(极验)v3 滑块。这带来两个影响。

一是 drain_timeout 按最慢那类取值,按平均值设会稳定切掉尾部任务。二是 reCAPTCHA 依赖 Google 托管脚本,国内网络下并不总是可达,这类偶发失败属于环境噪声而非版本缺陷——版本正确性用 getbalance 判定,站点连通性单独告警,不参与发版门禁。

容量上 CaptchaAI 按并发线程计费,套餐线程数固定(BASIC 每月 $15、5 个线程;ADVANCE 每月 $90、50 个线程),套餐内识别次数不限。

常见问题

已提交但没取回结果的任务会丢吗?

不会,前提是排空逻辑正确。captcha_id 已在服务端排队,节点只需把轮询走完,drain_timeout 要覆盖最长一次轮询周期。

健康检查会额外消耗额度吗?

日常发版用 getbalance 就够,只做一次鉴权往返。核心链路才值得跑真实识别任务——按线程计费,代价只是一小段线程占用。

回滚后版本不一致怎么办?

根因基本是回滚没覆盖全部已更新节点。维护 updated 列表,按列表逐个回退,再读各节点上报的 version 核对。

什么时候该换成蓝绿部署?

当逐个回退的时间超出故障容忍窗口时。滚动更新的回滚是线性的,节点越多越慢;蓝绿把回滚变成一次流量切换,代价是维持两套环境。

下一步

把发版升级成可排空、可回滚的滚动更新——获取你的 CaptchaAI API Key,用上面的协调器跑通一轮。

相关指南:

该文章已禁用评论。