DevOps & Scaling

用于大规模验证码解决的 Kubernetes 作业队列

验证码识别的流量常常忽高忽低:平时几百个任务,遇到批量采集或压测就瞬间涨到几万。固定数量的 worker 不是闲置浪费,就是被积压压垮。把任务写进 Redis 队列、用 Kubernetes 部署 worker Pod,再让 HPA 按队列深度自动扩缩容,算力就能跟着积压量走。下面是一套可直接落地的清单。


架构概览

链路里只有四个角色:生产者写入队列,worker 取任务并调用 CaptchaAI API,结果写回 Redis,HPA 按积压量增减 worker。

Producer → Redis Queue → Worker Pods (auto-scaled) → CaptchaAI API
                              ↓
                       Results Store (Redis)

worker 部署清单

先起 3 个副本:API Key 用 Secret 注入,Redis 地址走环境变量,每个 Pod 都设好 CPU 与内存的 requests/limits。

# k8s/worker-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
  name: captcha-worker
  labels:
    app: captcha-worker
spec:
  replicas: 3
  selector:
    matchLabels:
      app: captcha-worker
  template:
    metadata:
      labels:
        app: captcha-worker
    spec:
      containers:

        - name: worker
          image: your-registry/captcha-worker:latest
          env:

            - name: CAPTCHAAI_KEY
              valueFrom:
                secretKeyRef:
                  name: captchaai-secret
                  key: api-key

            - name: REDIS_URL
              value: "redis://redis-service:6379"
          resources:
            requests:
              memory: "128Mi"
              cpu: "100m"
            limits:
              memory: "256Mi"
              cpu: "250m"

用 Secret 管理 API Key

密钥不要写进镜像或 YAML,用 Secret 单独保管:

kubectl create secret generic captchaai-secret \
  --from-literal=api-key=YOUR_API_KEY

部署 Redis 队列

队列和结果都存在 Redis,下面同时创建 Deployment 和 Service,worker 通过 redis-service 访问:

# k8s/redis.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
  name: redis
spec:
  replicas: 1
  selector:
    matchLabels:
      app: redis
  template:
    metadata:
      labels:
        app: redis
    spec:
      containers:

        - name: redis
          image: redis:7-alpine
          ports:

            - containerPort: 6379
          resources:
            requests:
              memory: "128Mi"
              cpu: "100m"
---
apiVersion: v1
kind: Service
metadata:
  name: redis-service
spec:
  selector:
    app: redis
  ports:

    - port: 6379

worker 处理逻辑

主循环用 blpop 取任务,调用 in.php 提交、每 5 秒轮询一次 res.php,拿到 token 后写回结果,并刷新队列长度指标供 HPA 使用。

# worker.py
import os
import json
import time
import redis
import requests


class CaptchaWorker:
    """Kubernetes worker that processes CAPTCHA tasks from Redis."""

    def __init__(self):
        self.api_key = os.environ["CAPTCHAAI_KEY"]
        self.redis = redis.from_url(
            os.environ.get("REDIS_URL", "redis://localhost:6379"),
        )
        self.base = "https://ocr.captchaai.com"

    def run(self):
        """Main worker loop."""
        hostname = os.environ.get("HOSTNAME", "unknown")
        print(f"Worker {hostname} started")

        while True:
            result = self.redis.blpop("captcha:queue", timeout=30)
            if result is None:
                continue

            _, raw = result
            task = json.loads(raw)
            task_id = task.get("id", "unknown")

            print(f"[{hostname}] Processing {task_id}")
            start = time.time()

            try:
                token = self._solve(task["method"], task["params"])
                duration = time.time() - start
                self.redis.hset("captcha:results", task_id, json.dumps({
                    "status": "success",
                    "token": token,
                    "duration": f"{duration:.1f}s",
                    "worker": hostname,
                }))
                print(f"[{hostname}] {task_id} solved in {duration:.1f}s")

            except Exception as e:
                self.redis.hset("captcha:results", task_id, json.dumps({
                    "status": "error",
                    "error": str(e),
                    "worker": hostname,
                }))
                print(f"[{hostname}] {task_id} failed: {e}")

            # Update queue length metric
            queue_len = self.redis.llen("captcha:queue")
            self.redis.set("captcha:queue_length", queue_len)

    def _solve(self, method, params, timeout=120):
        resp = requests.post(f"{self.base}/in.php", data={
            "key": self.api_key,
            "method": method,
            "json": 1,
            **params,
        }, 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")


if __name__ == "__main__":
    CaptchaWorker().run()

用 HPA 按队列深度自动扩缩容

把触发器绑到队列长度而不是 CPU:积压越多,worker 越多。下例 2 到 20 个副本,平均每 10 个排队任务扩出一个 Pod。

# k8s/hpa.yaml
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
  name: captcha-worker-hpa
spec:
  scaleTargetRef:
    apiVersion: apps/v1
    kind: Deployment
    name: captcha-worker
  minReplicas: 2
  maxReplicas: 20
  metrics:

    - type: External
      external:
        metric:
          name: redis_queue_length
          selector:
            matchLabels:
              queue: captcha
        target:
          type: AverageValue
          averageValue: "10"

任务生产者:向队列投递任务

给每个任务生成短 ID、rpush 进队列,再按 ID 轮询结果:

import json
import uuid
import redis


def submit_tasks(redis_url, tasks):
    """Submit CAPTCHA tasks to the queue."""
    r = redis.from_url(redis_url)
    task_ids = []

    for task in tasks:
        task_id = str(uuid.uuid4())[:8]
        task["id"] = task_id
        r.rpush("captcha:queue", json.dumps(task))
        task_ids.append(task_id)

    return task_ids


def get_results(redis_url, task_ids, timeout=180):
    """Wait for and collect results."""
    r = redis.from_url(redis_url)
    results = {}
    deadline = time.time() + timeout

    while len(results) < len(task_ids) and time.time() < deadline:
        for tid in task_ids:
            if tid in results:
                continue
            raw = r.hget("captcha:results", tid)
            if raw:
                results[tid] = json.loads(raw)
        time.sleep(1)

    return results

把 worker 并发和 CaptchaAI 线程数对齐

CaptchaAI 按并发线程计费,每个正在处理的验证码占用一个线程。如果 HPA 把 worker 扩到上限、每个 Pod 又跑多个并发,总并发一旦超过套餐的线程额度,多出来的请求只会在 API 侧排队,扩容也没有收益。按套餐倒推上限最稳妥:ADVANCE($90/月,50 线程)约支撑 50 路并发,PREMIUM($170/月,100 线程)支撑 100 路。做数据采集的团队常在大促高峰临时放大额度、平峰再缩回——采集范围请遵循 robots 协议与个人信息保护法(PIPL)。


排错速查

现象 可能原因 处理方式
worker 起不来 没创建 Secret 先执行 kubectl create secret
Pod 反复 CrashLoopBackOff 缺环境变量或连不上 Redis kubectl logs 看日志
HPA 不扩容 没接外部/自定义指标 装指标适配器或改用 KEDA
队列只涨不消 worker 空转或崩溃 检查 Pod 状态并重启

常见问题

worker 并发怎么和 CaptchaAI 线程数匹配?

把「maxReplicas × 单 Pod 并发」压在套餐线程数以内,超出的请求只会在 API 侧排队。要更高并发就升级到线程更多的套餐。

队列一直在涨、结果却出不来怎么办?

kubectl get pods 看有没有 CrashLoopBackOff,再用 kubectl logs 查 Redis 连接或 API Key 是否出错。多半是 worker 全部空转或已崩溃。

大陆访问 reCAPTCHA 脚本不稳,会影响这套架构吗?

不影响。识别在 CaptchaAI 侧完成,worker 只和 CaptchaAI API 通信,与本地页面能否加载 Google 脚本无关。

高峰过后怎么自动缩容省成本?

HPA 会在队列变短后自动缩容,把 minReplicas 设为 2 保底即可;想缩到 0 或按事件触发就改用 KEDA。


相关阅读


把识别能力扩展到数千并发——注册 CaptchaAI 驱动你的 Kubernetes 集群。

该文章已禁用评论。