DevOps & Scaling

自动缩放验证码解决工作人员

双 11、618 大促期间,验证请求量能在几分钟内翻好几倍。固定 worker 池要么闲时烧钱,要么高峰期超时排队——答案是让 worker 数量跟着队列深度、利用率、余额自动升降。下面是可以直接抄的 Python 代码,附故障排查清单。

先选型:三种扩缩容方案怎么选

先想清楚一件事:worker 是在等 CaptchaAI 响应,还是在本地做图像处理?前者是 I/O 密集型,后者是 CPU 密集型,答案决定用哪种方案:

  1. 纯 API 调用,没有本地重活 → 线程池:延迟低、实现简单,多数场景直接够用。
  2. 本地还带图像预处理、模型推理这类活 → 进程池:延迟中等,靠独立进程隔离 CPU 占用。
  3. 已经跑在 Kubernetes 上,要云原生弹性 → Kubernetes HPA:延迟较高、复杂度高,适合大规模部署。
  4. 想按队列长度等事件触发扩容 → 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,直接把余额熔断和队列监控接进你的脚本。

该文章已禁用评论。