文章
合集Python 并发编程第 6 / 7 篇

实战:LLM 并发调用

综合案例,用一个真实项目把前 5 篇知识串起来。我们构建一个批量调用 LLM 审核数据标注质量的系统,涵盖选型决策、并发实现、限速控制、错误处理。

1. 项目背景

假设你在做一个数据标注平台,需要审核 500 条数据标注的质量。每条数据要调用 LLM(如 OpenAI GPT-4)来判断标注是否正确。

需求:

  1. 批量发送 500 个请求到 OpenAI API
  2. 每个请求需要审核一条数据标注
  3. API 有频率限制(每分钟 500 次)
  4. 需要统计成功/失败/超时的请求数
  5. 结果保持和输入相同的顺序

2. 选型决策:为什么选多线程

回到 concurrency-overview-selection 的决策树:

任务是 I/O 密集还是 CPU 密集?
│
├─ I/O 密集(等 API 响应,CPU 大部分时间空闲)
│   ├─ 任务数 500(< 100?不,> 100)
│   └─ 但 asyncio 需要改用异步库...

决策过程:

  • 任务是 I/O 密集(等网络响应):确定
  • OpenAI SDK 的 openai 库支持同步调用:用多线程最简单
  • 如果用 asyncio,需要改用 httpx.AsyncClient 或等 OpenAI SDK 支持异步(已支持)
  • 但当前项目已有大量同步代码,改 asyncio 成本太高

3. OpenAI SDK 调用模式

3.1 基本调用

from openai import OpenAI

client = OpenAI(api_key="sk-...")

def review_annotation(data_item: dict) -> dict:
    """调用 LLM 审核单条数据标注。"""
    response = client.chat.completions.create(
        model="gpt-4",
        messages=[
            {"role": "system", "content": "你是数据标注审核专家。"},
            {"role": "user", "content": f"审核以下标注:{data_item}"},
        ],
        temperature=0.1,   # 低温度,结果更确定
        response_format={"type": "json_object"},   # 要求 JSON 输出
    )
    return {
        "item_id": data_item["id"],
        "result": response.choices[0].message.content,
        "tokens_used": response.usage.total_tokens,
    }

3.2 关键参数

参数作用建议
model使用的模型根据任务复杂度选择
temperature输出随机性审核任务用 0.1,创意任务用 0.7+
response_format输出格式审核任务用 json_object,便于解析
timeout请求超时SDK 默认 600 秒,建议缩短

4. ThreadPoolExecutor 实现并发

4.1 核心实现

import logging
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from openai import OpenAI

logger = logging.getLogger(__name__)

def batch_review(data_items: list[dict], max_workers: int = 10) -> list[dict]:
    """并发批量审核数据标注。"""
    client = OpenAI(api_key="sk-...")   # 每个线程一个 client(线程安全)

    results = []
    with ThreadPoolExecutor(max_workers=max_workers) as executor:
        # 提交所有任务
        futures = {
            executor.submit(review_annotation, item): item
            for item in data_items
        }

        # 按完成顺序处理结果
        for future in as_completed(futures):
            item = futures[future]
            try:
                result = future.result(timeout=30)   # 单个任务超时
                results.append(result)
                logger.info("审核完成: %s", item["id"])
            except Exception as e:
                logger.error("审核失败 [%s]: %s", item["id"], e)
                results.append({
                    "item_id": item["id"],
                    "result": None,
                    "error": str(e),
                })

    return results

4.2 设计要点

  1. as_completed:先完成的先处理,不阻塞
  2. 每个 future 单独 try/except:一个失败不影响其他的
  3. future.result(timeout=30):防止单个任务卡死
  4. 错误记录包含 item_id:方便排查问题

5. Lock 限速器

API 有频率限制(如每分钟 500 次请求),需要限速:

import threading
import time

class RateLimiter:
    """基于时间窗口的限速器。"""

    def __init__(self, max_requests: int, window_seconds: int):
        self.max_requests = max_requests
        self.window_seconds = window_seconds
        self.requests = []
        self.lock = threading.Lock()

    def wait_if_needed(self):
        """如果超过速率限制,等待到窗口重置。"""
        with self.lock:
            now = time.monotonic()

            # 移除窗口外的旧请求记录
            self.requests = [
                t for t in self.requests
                if now - t < self.window_seconds
            ]

            if len(self.requests) >= self.max_requests:
                # 窗口满了,等待最早的请求过期
                sleep_time = self.window_seconds - (now - self.requests[0])
                if sleep_time > 0:
                    logger.info("限速:等待 %.1f 秒", sleep_time)
                    time.sleep(sleep_time)
                    # 重新清理窗口
                    now = time.monotonic()
                    self.requests = [
                        t for t in self.requests
                        if now - t < self.window_seconds
                    ]

            self.requests.append(time.monotonic())

使用限速器

rate_limiter = RateLimiter(max_requests=50, window_seconds=60)   # 每分钟 50 次

def review_with_rate_limit(client, item):
    """带限速的审核。"""
    rate_limiter.wait_if_needed()   # 超过速率就等待
    return review_annotation(client, item)

def batch_review(data_items: list[dict], max_workers: int = 10) -> list[dict]:
    """带限速的并发审核。"""
    client = OpenAI(api_key="sk-...")

    with ThreadPoolExecutor(max_workers=max_workers) as executor:
        futures = {
            executor.submit(review_with_rate_limit, client, item): item
            for item in data_items
        }

        results = []
        for future in as_completed(futures):
            item = futures[future]
            try:
                result = future.result(timeout=30)
                results.append(result)
            except Exception as e:
                logger.error("审核失败 [%s]: %s", item["id"], e)
                results.append({
                    "item_id": item["id"],
                    "result": None,
                    "error": str(e),
                })

        return results

6. 结果收集与统计

6.1 保持顺序

as_completed 按完成顺序返回,如果需要保持输入顺序:

from concurrent.futures import ThreadPoolExecutor, as_completed

def batch_review_ordered(data_items: list[dict], max_workers: int = 10) -> list[dict]:
    """保持输入顺序的并发审核。"""
    client = OpenAI(api_key="sk-...")

    # 用 dict 记录每个 future 对应的索引
    indexed_futures = {}
    with ThreadPoolExecutor(max_workers=max_workers) as executor:
        for idx, item in enumerate(data_items):
            future = executor.submit(review_annotation, client, item)
            indexed_futures[future] = idx

        # 收集结果
        results = [None] * len(data_items)
        for future in as_completed(indexed_futures):
            idx = indexed_futures[future]
            try:
                results[idx] = future.result(timeout=30)
            except Exception as e:
                results[idx] = {"item_id": data_items[idx]["id"], "error": str(e)}

    return results   # 顺序和输入一致

6.2 统计信息

def print_stats(results: list[dict]):
    """打印审核统计。"""
    total = len(results)
    success = sum(1 for r in results if r.get("result") is not None)
    failed = sum(1 for r in results if "error" in r)
    tokens = sum(r.get("tokens_used", 0) for r in results)

    print(f"总计: {total} 条")
    print(f"成功: {success} 条")
    print(f"失败: {failed} 条")
    print(f"Token 消耗: {tokens}")

    if failed > 0:
        print("失败详情:")
        for r in results:
            if "error" in r:
                print(f"  - {r['item_id']}: {r['error']}")

7. 扩展讨论

7.1 如果改成 asyncio 怎么做

import asyncio
import httpx

async def review_annotation_async(client, item):
    """异步版本的审核。"""
    response = await client.post(
        "https://api.openai.com/v1/chat/completions",
        json={
            "model": "gpt-4",
            "messages": [
                {"role": "system", "content": "你是数据标注审核专家。"},
                {"role": "user", "content": f"审核以下标注:{item}"},
            ],
            "temperature": 0.1,
        },
        headers={"Authorization": f"Bearer {api_key}"},
    )
    return response.json()

async def batch_review_async(data_items, max_concurrent=50):
    """异步批量审核。"""
    semaphore = asyncio.Semaphore(max_concurrent)
    async with httpx.AsyncClient(timeout=30) as client:

        async def limited_review(item):
            async with semaphore:
                return await review_annotation_async(client, item)

        tasks = [limited_review(item) for item in data_items]
        return await asyncio.gather(*tasks, return_exceptions=True)

7.2 如果要多进程怎么办

LLM 调用是 I/O 密集,不需要多进程。但如果审核逻辑包含本地数据处理(CPU 密集),可以混合使用:

from concurrent.futures import ProcessPoolExecutor, ThreadPoolExecutor

# CPU 密集:本地数据预处理
with ProcessPoolExecutor() as executor:
    processed_data = list(executor.map(preprocess, raw_data))

# I/O 密集:LLM 调用
with ThreadPoolExecutor(max_workers=10) as executor:
    results = list(executor.map(review_annotation, processed_data))

8. 完整代码

"""
llm_review.py - 并发调用 LLM 审核数据标注质量
综合演示:ThreadPoolExecutor + Lock 限速器 + 错误处理
"""
import logging
import time
import threading
from concurrent.futures import ThreadPoolExecutor, as_completed
from openai import OpenAI

logger = logging.getLogger(__name__)


class RateLimiter:
    """基于时间窗口的限速器。"""

    def __init__(self, max_requests: int, window_seconds: int):
        self.max_requests = max_requests
        self.window_seconds = window_seconds
        self.requests: list[float] = []
        self.lock = threading.Lock()

    def wait_if_needed(self):
        with self.lock:
            now = time.monotonic()
            self.requests = [
                t for t in self.requests
                if now - t < self.window_seconds
            ]
            if len(self.requests) >= self.max_requests:
                sleep_time = self.window_seconds - (now - self.requests[0])
                if sleep_time > 0:
                    time.sleep(sleep_time)
                    now = time.monotonic()
                    self.requests = [
                        t for t in self.requests
                        if now - t < self.window_seconds
                    ]
            self.requests.append(time.monotonic())


def review_annotation(client: OpenAI, item: dict) -> dict:
    """调用 LLM 审核单条数据标注。"""
    response = client.chat.completions.create(
        model="gpt-4",
        messages=[
            {"role": "system", "content": "你是数据标注审核专家。"},
            {"role": "user", "content": f"审核以下标注:{item}"},
        ],
        temperature=0.1,
        response_format={"type": "json_object"},
    )
    return {
        "item_id": item["id"],
        "result": response.choices[0].message.content,
        "tokens_used": response.usage.total_tokens,
    }


def batch_review(
    data_items: list[dict],
    max_workers: int = 10,
    rate_limit: int = 50,
) -> list[dict]:
    """
    并发批量审核数据标注。

    Args:
        data_items: 待审核的数据列表
        max_workers: 并发线程数
        rate_limit: 每分钟最大请求数
    """
    client = OpenAI(api_key="sk-...")
    limiter = RateLimiter(max_requests=rate_limit, window_seconds=60)

    results: list[dict] = []
    with ThreadPoolExecutor(max_workers=max_workers) as executor:
        futures = {}
        for item in data_items:
            future = executor.submit(
                lambda item: (
                    limiter.wait_if_needed() or review_annotation(client, item)
                ),
                item,
            )
            futures[future] = item

        for future in as_completed(futures):
            item = futures[future]
            try:
                result = future.result(timeout=30)
                results.append(result)
            except Exception as e:
                logger.error("审核失败 [%s]: %s", item["id"], e)
                results.append({
                    "item_id": item["id"],
                    "result": None,
                    "error": str(e),
                })

    # 打印统计
    total = len(results)
    success = sum(1 for r in results if r.get("result") is not None)
    failed = sum(1 for r in results if "error" in r)
    logger.info("审核完成: %d/%d 成功, %d 失败", success, total, failed)

    return results


if __name__ == "__main__":
    logging.basicConfig(level=logging.INFO)

    # 模拟 500 条待审核数据
    data_items = [
        {"id": i, "text": f"样本 {i}", "label": "positive"}
        for i in range(500)
    ]

    start = time.perf_counter()
    results = batch_review(data_items, max_workers=10, rate_limit=50)
    elapsed = time.perf_counter() - start

    print(f"完成 {len(results)} 条审核,耗时 {elapsed:.1f} 秒")

9. 知识点回顾

这篇文章用到的知识:

知识点来源在本项目中的应用
I/O 密集判断concurrency-overview-selection决定用多线程
ThreadPoolExecutorpython-multithreading并发执行 LLM 调用
as_completedpython-multithreading按完成顺序处理结果
Lock 限速器python-multithreading控制 API 调用频率
错误隔离python-multithreading每个 future 单独 try/except
asyncio 对比python-asyncio讨论了改异步的方案
subprocess 对比python-subprocess如果要调外部脚本

10. 小结

这个实战项目用到了几个并发编程的核心模式:

  1. 选型:I/O 密集 + 同步 SDK → 多线程
  2. 实现:ThreadPoolExecutor + as_completed
  3. 限速:Lock + 时间窗口实现 Rate Limiter
  4. 容错:每个任务独立 try/except,一个失败不影响其他
  5. 统计:收集结果,打印成功/失败/Token 消耗