实战:LLM 并发调用
综合案例,用一个真实项目把前 5 篇知识串起来。我们构建一个批量调用 LLM 审核数据标注质量的系统,涵盖选型决策、并发实现、限速控制、错误处理。
1. 项目背景
假设你在做一个数据标注平台,需要审核 500 条数据标注的质量。每条数据要调用 LLM(如 OpenAI GPT-4)来判断标注是否正确。
需求:
- 批量发送 500 个请求到 OpenAI API
- 每个请求需要审核一条数据标注
- API 有频率限制(每分钟 500 次)
- 需要统计成功/失败/超时的请求数
- 结果保持和输入相同的顺序
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 设计要点
as_completed:先完成的先处理,不阻塞- 每个 future 单独 try/except:一个失败不影响其他的
future.result(timeout=30):防止单个任务卡死- 错误记录包含 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 | 决定用多线程 |
ThreadPoolExecutor | python-multithreading | 并发执行 LLM 调用 |
as_completed | python-multithreading | 按完成顺序处理结果 |
Lock 限速器 | python-multithreading | 控制 API 调用频率 |
| 错误隔离 | python-multithreading | 每个 future 单独 try/except |
| asyncio 对比 | python-asyncio | 讨论了改异步的方案 |
| subprocess 对比 | python-subprocess | 如果要调外部脚本 |
10. 小结
这个实战项目用到了几个并发编程的核心模式:
- 选型:I/O 密集 + 同步 SDK → 多线程
- 实现:
ThreadPoolExecutor+as_completed - 限速:
Lock+ 时间窗口实现 Rate Limiter - 容错:每个任务独立 try/except,一个失败不影响其他
- 统计:收集结果,打印成功/失败/Token 消耗