双 11、618 大促期间,验证请求量能在几分钟内翻好几倍。固定 worker 池要么闲时烧钱,要么高峰期超时排队——答案是让 worker 数量跟着队列深度、利用率、余额自动升降。下面是可以直接抄的 Python 代码,附故障排查清单。
先选型:三种扩缩容方案怎么选
先想清楚一件事:worker 是在等 CaptchaAI 响应,还是在本地做图像处理?前者是 I/O 密集型,后者是 CPU 密集型,答案决定用哪种方案:
- 纯 API 调用,没有本地重活 → 线程池:延迟低、实现简单,多数场景直接够用。
- 本地还带图像预处理、模型推理这类活 → 进程池:延迟中等,靠独立进程隔离 CPU 占用。
- 已经跑在 Kubernetes 上,要云原生弹性 → Kubernetes HPA:延迟较高、复杂度高,适合大规模部署。
- 想按队列长度等事件触发扩容 → KEDA:延迟中等、复杂度中,适合事件驱动架构。
扩容信号:这四个指标说了算
选好方案后,扩缩容脚本该盯紧下面几个信号,而不是凭感觉调整数量:
- 队列深度:待处理任务 > 20 时扩容,< 5 时缩容。
- worker 利用率:忙碌 > 80% 时扩容,< 20% 时缩容。
- 识别耗时(P95):超过 60 秒扩容,低于 20 秒缩容。
- 错误率:超过 5% 说明该换新 worker 了,稳定低于 1% 才算健康。
- 余额:低于 $1 直接停止扩容,这是硬性熔断线。
队列深度和利用率最先该看,多数场景下够用;P95 冲到 60 秒以上就说明 worker 跟不上了。
线程池扩缩容器:单进程内动态加人
worker 维护在单个 Python 进程里,靠 Redis 队列分发任务:一个循环处理任务,一个扩缩容循环每 10 秒检查队列深度和利用率。
import os
import time
import threading
import requests
import json
import redis
class AutoScalingPool:
"""Dynamically scale CaptchaAI worker threads."""
def __init__(self, api_key, redis_url="redis://localhost:6379"):
self.api_key = api_key
self.redis = redis.from_url(redis_url)
self.base = "https://ocr.captchaai.com"
self.queue_key = "captcha:tasks"
self.results_key = "captcha:results"
self.min_workers = 2
self.max_workers = 20
self.workers = []
self.active_count = 0
self.lock = threading.Lock()
self.running = True
def start(self):
"""Start the pool with minimum workers."""
for _ in range(self.min_workers):
self._add_worker()
# Start scaler in background
scaler = threading.Thread(target=self._scaling_loop, daemon=True)
scaler.start()
print(f"Pool started with {self.min_workers} workers")
def _add_worker(self):
"""Add a worker thread."""
if len(self.workers) >= self.max_workers:
return
t = threading.Thread(target=self._worker_loop, daemon=True)
t.start()
self.workers.append(t)
def _remove_worker(self):
"""Signal one worker to stop (lazy removal)."""
if len(self.workers) <= self.min_workers:
return
self.workers.pop() # Thread will exit on next idle cycle
def _worker_loop(self):
"""Worker loop: fetch and process tasks."""
while self.running and threading.current_thread() in self.workers:
result = self.redis.blpop(self.queue_key, timeout=10)
if result is None:
continue
_, raw = result
task = json.loads(raw)
task_id = task["id"]
with self.lock:
self.active_count += 1
try:
token = self._solve(task["method"], task["params"])
self.redis.hset(self.results_key, task_id, json.dumps({
"status": "success", "token": token,
}))
except Exception as e:
self.redis.hset(self.results_key, task_id, json.dumps({
"status": "error", "error": str(e),
}))
finally:
with self.lock:
self.active_count -= 1
def _scaling_loop(self):
"""Periodically adjust worker count."""
while self.running:
time.sleep(10)
queue_depth = self.redis.llen(self.queue_key)
current = len(self.workers)
utilization = (
self.active_count / current * 100 if current > 0 else 0
)
# Scale up: queue growing and workers busy
if queue_depth > 20 and utilization > 70:
new_count = min(current + 2, self.max_workers)
while len(self.workers) < new_count:
self._add_worker()
print(f"Scaled up: {current} → {len(self.workers)} workers")
# Scale down: queue empty and workers idle
elif queue_depth < 5 and utilization < 20:
target = max(current - 1, self.min_workers)
while len(self.workers) > target:
self._remove_worker()
if len(self.workers) < current:
print(f"Scaled down: {current} → {len(self.workers)} workers")
def _solve(self, method, params, timeout=120):
data = {"key": self.api_key, "method": method, "json": 1}
data.update(params)
resp = requests.post(
f"{self.base}/in.php", data=data, timeout=30,
)
result = resp.json()
if result.get("status") != 1:
raise RuntimeError(result.get("request"))
captcha_id = result["request"]
start = time.time()
while time.time() - start < timeout:
time.sleep(5)
resp = requests.get(f"{self.base}/res.php", params={
"key": self.api_key,
"action": "get",
"id": captcha_id,
"json": 1,
}, timeout=15)
data = resp.json()
if data["request"] != "CAPCHA_NOT_READY":
if data.get("status") == 1:
return data["request"]
raise RuntimeError(data["request"])
raise TimeoutError("Solve timeout")
def stats(self):
return {
"workers": len(self.workers),
"active": self.active_count,
"queue": self.redis.llen(self.queue_key),
}
# Usage
pool = AutoScalingPool(os.environ["CAPTCHAAI_KEY"])
pool.start()
# Monitor
while True:
print(pool.stats())
time.sleep(30)
_scaling_loop 每 10 秒跑一次:队列超过 20 且利用率超过 70% 就加两个 worker;队列低于 5 且利用率低于 20% 才缩一个,缩容比扩容更保守,避免抖动。
进程池扩缩容器:CPU 密集型任务给隔离
worker 里混了图像预处理、本地推理这类吃 CPU 的活时,进程池更合适——每个 worker 独立进程,天然隔离 CPU 占用。
import multiprocessing
import time
import redis
import os
class ProcessScaler:
"""Scale worker processes based on queue depth."""
def __init__(self, worker_fn, redis_url="redis://localhost:6379"):
self.worker_fn = worker_fn
self.redis = redis.from_url(redis_url)
self.processes = []
self.min_workers = 2
self.max_workers = 16
def run(self, check_interval=15):
"""Run the scaler loop."""
# Start minimum workers
for _ in range(self.min_workers):
self._spawn()
while True:
time.sleep(check_interval)
self._cleanup_dead()
queue_depth = self.redis.llen("captcha:tasks")
current = len(self.processes)
# Scale up
if queue_depth > current * 5 and current < self.max_workers:
to_add = min(
max(1, queue_depth // 10),
self.max_workers - current,
)
for _ in range(to_add):
self._spawn()
print(f"Scaled up to {len(self.processes)} workers")
# Scale down
elif queue_depth < 3 and current > self.min_workers:
to_remove = min(2, current - self.min_workers)
for _ in range(to_remove):
p = self.processes.pop()
p.terminate()
print(f"Scaled down to {len(self.processes)} workers")
def _spawn(self):
p = multiprocessing.Process(target=self.worker_fn)
p.start()
self.processes.append(p)
def _cleanup_dead(self):
self.processes = [p for p in self.processes if p.is_alive()]
# Ensure minimum
while len(self.processes) < self.min_workers:
self._spawn()
ProcessScaler 按队列深度算需要新增的进程数,一次最多加到 max_workers;缩容每轮最多砍两个,_cleanup_dead() 顺手清掉僵尸进程。
余额感知扩容:别让脚本把余额刷穿
识别量上去了,余额消耗也跟着涨。上线前加一道余额检查:低于阈值就暂停扩容,而不是等余额耗尽后收到一堆失败请求。
def check_balance(api_key, min_balance=2.0):
"""Check if balance is sufficient for scaling."""
resp = requests.get("https://ocr.captchaai.com/res.php", params={
"key": api_key,
"action": "getbalance",
"json": 1,
}, timeout=15)
balance = float(resp.json()["request"])
if balance < min_balance:
print(f"Balance ${balance:.2f} below ${min_balance} — halting scale-up")
return False
return True
集成到缩放循环中:
# In _scaling_loop:
if queue_depth > 20 and utilization > 70:
if check_balance(self.api_key, min_balance=2.0):
# Scale up
...
else:
print("Scaling paused — low balance")
提示: 把
check_balance()塞进扩容分支——余额够用才加 worker,不够就打日志、跳过这一轮,账户充值后下一次检查自动恢复。
常见故障排查
- worker 数量一直往上涨:队列根本没排空 → 检查 worker 是否真的在处理任务,而不是卡住了。
- 缩容太激进,来回抖动:缩容阈值触发太快 → 把缩容延迟拉长到 30 秒以上。
- 出现僵尸进程:进程没被正确清理 → 定期调用
_cleanup_dead()。 - 余额掉得特别快:worker 开太多了 → 在扩容逻辑里加余额检查。
常见问题
worker 和队列的比例怎么定最合适?
经验值:每 5–10 个排队任务配 1 个 worker,每个 worker 每分钟约处理 3–6 个任务。
出现僵尸进程或者 worker 卡住不退出,怎么排查?
僵尸进程通常是子进程崩溃后没清理,定期跑 _cleanup_dead() 能自动摘掉;worker 只涨不跌,多半是队列没排空。
余额不够了,怎么让扩容自动停下来而不是报错?
在扩容分支加一次 check_balance():余额不足就跳过这轮并记日志,不会被失败请求刷屏。
到底该用线程池还是进程池?
纯调 API 用线程池就够;夹了图像预处理或本地推理这类吃 CPU 的步骤,才需要换成进程池。
相关指南
扩容这件事没必要靠猜——注册 CaptchaAI 拿到 API Key,直接把余额熔断和队列监控接进你的脚本。