升级一组验证码解决人员不应放弃正在进行的任务。滚动更新一次替换一个工作人员——耗尽活动任务,部署新版本,并在转移到下一个工作人员之前验证运行状况。
滚动更新流程
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
Python——滚动更新协调器
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}")
JavaScript – 带进度跟踪的滚动更新
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)));
更新策略比较
| 战略 | 停机时间 | 回滚速度 | 复杂 | 最适合 |
|---|---|---|---|---|
| 滚动 | 没有任何 | 缓和 | 低的 | 大多数部署 |
| 蓝绿 | 没有任何 | 立即的 | 中等的 | 关键服务 |
| 金丝雀 | 没有任何 | 快速地 | 高的 | 大型车队(50 名以上工人) |
| 重新创造 | 简短的 | N/A | 最低 | 开发环境 |
故障排除
| 问题 | 原因 | 处理方式 |
|---|---|---|
| 更新期间任务被丢弃 | 排水超时太短 | 增加drain_timeout;检查最大任务持续时间 |
| 更新后健康检查总是失败 | 新版本有bug | 回滚;首先在暂存中测试新版本 |
| 更新时间太长 | maxUnavailable 太低,有很多工人 |
对于大型车队,将 maxUnavailable 增加到 2-3 |
| 回滚留下混合版本 | 部分回滚不完整 | 跟踪更新了哪些工作人员;全部回滚 |
常问问题
什么是好的 maxUnavailable 设置?
对于小型车队(< 10 名工人),一次更新 1 个。对于较大的车队,使用员工总数的 10-25%。切勿超过 50% - 您需要足够的容量来处理交通。
排水超时应该多长时间?
将其设置为最大验证码解决时间加上缓冲区。对于 reCAPTCHA v2(通常为 30-90 秒),120 秒的耗尽超时是安全的。
我应该在每次工作人员更新后运行集成测试吗?
对于关键系统是的。在移动到下一个之前,对每个新更新的工作人员运行真正的验证码解决方案。这可以及早捕获配置错误。
下一步
无需停机即可部署更新 -获取您的 CaptchaAI API 密钥并实施滚动更新。
相关指南: