给验证码识别 worker 集群发版,最稳的做法是:一次只动一个节点,先排空在途任务,再换版本,健康检查通过才轮到下一个。全程都有容量在跑,调用方无感。
直接全部重启的代价很实在:验证码识别是长耗时异步任务,提交到 in.php 拿 captcha_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,用上面的协调器跑通一轮。
相关指南: