Tutorials

用于验证码解决并行性的 Python ThreadPoolExecutor

验证码识别慢,很多时候不是 CaptchaAI 的问题,而是脚本在傻等 HTTP 响应。顺序请求 20 个验证码,等待时间会直接叠加成几分钟;Python 内置的 ThreadPoolExecutor 几行代码就能让等待并行,不用把调用链改写成 asyncio

本文覆盖:

  • 批量提交、轮询、Session 复用、超时与进度控制
  • max_workers 怎么选,以及和 asyncio 的取舍

为什么验证码识别适合用 ThreadPoolExecutor 并行

验证码识别是 I/O-bound 任务,大部分时间在等 HTTP 响应;Python 线程在 I/O 期间释放 GIL,ThreadPoolExecutor 因此高效:

  • 顺序执行:零改动,无并行能力
  • ThreadPoolExecutor:改动低,I/O 并行能力好
  • asyncio:并行能力最佳,但要重写成异步
  • multiprocessing:兼容现有代码,对 I/O 场景是过度设计

常见组合:

  • 国内站点:GeeTest(极验)
  • 出海站点:reCAPTCHA 或 Turnstile

ThreadPoolExecutor 对两类目标都能并行处理。

基础实现:批量提交与轮询

  • 提交一个 reCAPTCHA v2 任务
  • 轮询直到拿到结果,逻辑封装进一个同步函数
import os
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
import requests

API_KEY = os.environ["CAPTCHAAI_API_KEY"]


def solve_captcha(sitekey, pageurl):
    """Synchronous CAPTCHA solve — submit and poll."""
    # Submit
    resp = requests.post("https://ocr.captchaai.com/in.php", data={
        "key": API_KEY,
        "method": "userrecaptcha",
        "googlekey": sitekey,
        "pageurl": pageurl,
        "json": 1
    })
    data = resp.json()

    if data.get("status") != 1:
        raise RuntimeError(data.get("request", "Submit failed"))

    captcha_id = data["request"]

    # Poll for result
    for _ in range(60):
        time.sleep(5)
        result = requests.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 result["request"]
        if result.get("request") != "CAPCHA_NOT_READY":
            raise RuntimeError(result.get("request", "Unknown error"))

    raise TimeoutError("Solve timeout after 300s")


# Batch solve with ThreadPoolExecutor
tasks = [
    {"sitekey": "6Le-wvkSAAAAAPBMRTvw0Q4Muexq9bi0DJwx_mJ-", "pageurl": f"https://example.com/page/{i}"}
    for i in range(20)
]

start = time.time()

with ThreadPoolExecutor(max_workers=10) as executor:
    futures = {
        executor.submit(solve_captcha, t["sitekey"], t["pageurl"]): t
        for t in tasks
    }

    solved = 0
    failed = 0

    for future in as_completed(futures):
        task = futures[future]
        try:
            solution = future.result()
            solved += 1
            print(f"[OK] {task['pageurl']}: {solution[:30]}...")
        except Exception as e:
            failed += 1
            print(f"[ERR] {task['pageurl']}: {e}")

elapsed = time.time() - start
print(f"\nDone: {solved} solved, {failed} failed in {elapsed:.1f}s")

用 Session 复用连接,降低握手开销

每次请求都新建一个 TCP 连接会浪费时间。让每个线程复用同一个 requests.Session

import threading

# Thread-local storage for sessions
thread_local = threading.local()


def get_session():
    """Get or create a thread-local session."""
    if not hasattr(thread_local, "session"):
        thread_local.session = requests.Session()
        # Configure connection pooling
        adapter = requests.adapters.HTTPAdapter(
            pool_connections=10,
            pool_maxsize=10,
            max_retries=2
        )
        thread_local.session.mount("https://", adapter)
    return thread_local.session


def solve_captcha_pooled(sitekey, pageurl):
    """Solve using thread-local connection pooling."""
    session = get_session()

    resp = session.post("https://ocr.captchaai.com/in.php", data={
        "key": API_KEY,
        "method": "userrecaptcha",
        "googlekey": sitekey,
        "pageurl": pageurl,
        "json": 1
    })
    data = resp.json()

    if data.get("status") != 1:
        raise RuntimeError(data.get("request"))

    captcha_id = data["request"]

    for _ in range(60):
        time.sleep(5)
        result = 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 result["request"]
        if result.get("request") != "CAPCHA_NOT_READY":
            raise RuntimeError(result.get("request"))

    raise TimeoutError("Solve timeout")

提示:pool_maxsize 建议和 max_workers 保持一致,否则连接池会成为新的瓶颈。

用 map() 做简单批量任务

  • 不需要逐个任务单独处理异常
  • 只关心批量结果列表,顺序无所谓
def solve_task(task):
    """Wrapper that returns result dict."""
    try:
        solution = solve_captcha_pooled(task["sitekey"], task["pageurl"])
        return {"url": task["pageurl"], "solution": solution, "error": None}
    except Exception as e:
        return {"url": task["pageurl"], "solution": None, "error": str(e)}


with ThreadPoolExecutor(max_workers=10) as executor:
    results = list(executor.map(solve_task, tasks))

solved = [r for r in results if r["solution"]]
failed = [r for r in results if r["error"]]
print(f"Solved: {len(solved)}, Failed: {len(failed)}")

超时保护,避免线程池被卡死

超时保护要解决两件事:

  • 单个任务卡住不该拖累整个批次
  • 全局超时和单任务超时分开设置
from concurrent.futures import TimeoutError as FuturesTimeout

with ThreadPoolExecutor(max_workers=10) as executor:
    futures = {
        executor.submit(solve_captcha_pooled, t["sitekey"], t["pageurl"]): t
        for t in tasks
    }

    for future in as_completed(futures, timeout=600):  # 10 min global timeout
        task = futures[future]
        try:
            solution = future.result(timeout=120)  # 2 min per task
            print(f"[OK] {task['pageurl']}")
        except FuturesTimeout:
            print(f"[TIMEOUT] {task['pageurl']}")
        except Exception as e:
            print(f"[ERR] {task['pageurl']}: {e}")

实时进度回调

批次一长,最好能实时看到进度,方便判断卡在哪一步:

import threading

progress_lock = threading.Lock()
progress = {"done": 0, "total": 0}


def solve_with_progress(task):
    result = solve_task(task)
    with progress_lock:
        progress["done"] += 1
        pct = progress["done"] / progress["total"] * 100
        print(f'\r  Progress: {progress["done"]}/{progress["total"]} ({pct:.0f}%)', end="")
    return result


progress["total"] = len(tasks)

with ThreadPoolExecutor(max_workers=10) as executor:
    results = list(executor.map(solve_with_progress, tasks))

print()  # Newline after progress

如何选择 max_workers

  • 5 线程:开销很低,小批量稳妥起步
  • 10 线程:开销低,常规场景
  • 25 线程:开销中等,高吞吐流水线
  • 50 线程:开销较高,追求最大吞吐

起步建议:

  • max_workers=10 开始,边观察错误率边往上调
  • 别超过套餐线程上限:BASIC $15/月 5 线程,ADVANCE $90/月 50 线程

ThreadPoolExecutor 该选它还是 asyncio?

# ThreadPoolExecutor — drop into existing sync code
with ThreadPoolExecutor(max_workers=10) as executor:
    results = list(executor.map(solve_task, tasks))

# asyncio — requires async function chain
async def main():
    async with aiohttp.ClientSession() as session:
        tasks = [solve_async(session, t) for t in task_list]
        results = await asyncio.gather(*tasks)

ThreadPoolExecutor:代码是同步的、依赖 Selenium 这类不支持 async 的库、想快速并行又不想大改架构。

asyncio:项目从零搭建、追求最少的系统线程开销、已经在用 FastAPI、aiohttp。

常见故障排查

现象 原因 处理方式
所有线程都卡住 轮询时被 time.sleep 挂起 正常——sleep 期间会释放 GIL
ConnectionError 增多 并发连接数太高 调低 max_workers;启用连接池
结果顺序乱了 as_completed 按完成顺序返回 map() 保序,或自己用 dict 记录
内存持续增长 future 堆积的大对象未释放 as_completed 循环里边拿边处理

提示:先调低 max_workers 复现问题,能显著缩小排查范围。

常见问题

GIL 会限制 ThreadPoolExecutor 的真正并行吗?

不会。I/O 等待(HTTP 请求、time.sleep)时 Python 会释放 GIL,线程能真正并发;GIL 只限制 CPU 密集型任务。

小项目该用 ThreadPoolExecutor 还是直接上 asyncio?

代码已是同步的、或依赖 Selenium,选 ThreadPoolExecutor 更省事;全新项目且并发量大,选 asyncio

用 ProcessPoolExecutor 会不会更快?

不会,反而更慢——ProcessPoolExecutor 只增加进程间通信开销,I/O-bound 场景该用线程。

Selenium 脚本能直接套用这套并行方案吗?

可以。把提交和等待逻辑包进一个函数,用线程池并行调用即可,不用改写成异步。

下一步

注册 CaptchaAI 拿 API Key,把 ThreadPoolExecutor 接进现有流程,几行代码就能把顺序等待变成并发请求。

相关指南: 并行验证码识别方案并行与顺序处理性能对比,以及每小时处理 10,000 个任务的实践

该文章已禁用评论。