DeepSeek批量请求与异步调用实战:把2000条文本处理从40分钟压缩到3分半
简介面向希望充分发挥 DeepSeek 性能的开发者与数据工作者本资源是一份以批量请求与异步调用为主线的实战型 PDF 文档共 21 页内容覆盖环境准备、API 密钥获取、批量请求构建与发送、响应解析、错误处理与重试机制以及基于 asyncio 和 async/await 的异步调用实现与并发控制方法并配有性能测试与优化章节帮助读者掌握从基础原理到项目落地的完整路径。文档还专门梳理了请求超时、部分请求失败、异步任务异常等常见问题及对应解决思路实用性强。文件总数 1 个为 PDF 格式压缩包大小 1.74MB内容排版与图表显示完整便于直接阅读与随时查阅。已有 63 人浏览学习适合正在使用 DeepSeek 进行数据处理、内容生成或高并发场景开发并希望缩短请求耗时、提升系统吞吐量的初中级技术人员。1. 把 2000 条文本跑完从 40 分钟压到 3 分半DeepSeek 批量请求与异步调用的真实收益我在一个真实项目里第一次意识到DeepSeek 调用方式的差别远比模型本身更影响交付时间。那会儿要一天跑完约 2000 条电商评论让模型做情感分类和标签抽取最开始是 for 循环逐个请求跑了 40 分钟还在继续。后来把请求组织成批量、再把等待时间交给事件循环同一批数据 3 分半跑完。这篇要拆的批量请求与异步调用就是这类“文件级文本处理”场景——评论分类、OCR 结果清洗、批量摘要、客服工单打标——最省钱也最直接的提效手段。适合需要大量调 API、又不想因为高并发把账户打爆的开发者。2. 批量请求合并通信开销先把单条 for 循环换成一次往返2.1 批量请求的工作流程与选型理由批量请求的核心思路很简单把多个子请求打包成一个请求体一次连接、一次往返服务器解析后返回组合结果。相比逐个请求它省掉的是连接建立、TCP 握手、请求头重复传输这些固定开销。DeepSeek 这类大模型 API 单次请求动辄几百毫秒里面网络往返占的时间比例不低把 10 条文本捆在一起发省下的往往是数百微秒到数毫秒的单条开销批量数量越大收益越明显。选型时还有一个容易被忽视的理由服务端吞吐。每条请求都要经历鉴权、路由、模型推理排队拆成 1000 个单条请求服务端要处理 1000 次上下文切换合并成 50 个批量请求排队压力和连接数都降一个量级。我一般把批量请求作为“并发改造之前的第一版”因为它改动最小只需要调整请求组装逻辑。需要说明的是这里讲的“服务器返回组合结果”依赖你的 API 是否支持数组入参。多数大模型 API 的标准聊天补全接口更常见的是单条 messages 请求此时批量要做的是“逻辑层合并”——把任务分流到若干批次每批只发一个请求而不是把一个数组直接塞给聊天补全接口。你拿到的 DeepSeek API 以哪种方式可用第一件事就是看文档里的请求体格式。2.2 构建批量数据、发送请求并解析响应先装依赖之后的所有示例都基于requestspip install requests然后写一个能直接跑通的批量处理骨架import requests import re import time API_KEY sk-你的密钥 API_URL https://api.deepseek.com/chat/completions raw_items [ {id: 1, text: 物流很快 包装有点破损}, {id: 2, text: 电池续航一般 充电速度倒是快}, {id: 3, text: 客服态度很好 退货也简单}, ] def preprocess_text(text): # 清理多余空白避免把空字符也当成输入 return re.sub(r\s, , text).strip() # 构造批量 payload每个元素是一个完整的聊天补全请求 batch_data [] for item in raw_items: batch_data.append({ model: deepseek-chat, # 以你账户可见的模型标识为准 messages: [ {role: system, content: 你是一个评论分类助手只输出结果。}, {role: user, content: preprocess_text(item[text])} ], temperature: 0.2, }) headers { Content-Type: application/json, Authorization: fBearer {API_KEY} } max_retries 3 retry_count 0 response None while retry_count max_retries: try: response requests.post(API_URL, headersheaders, jsonbatch_data, timeout30) response.raise_for_status() break except requests.exceptions.Timeout: print(f请求超时当前第 {retry_count 1} 次重试) except requests.exceptions.HTTPError as e: if e.response.status_code in (429, 500, 502, 503): retry_count 1 time.sleep(2 ** retry_count) else: print(fHTTP 错误{e.response.status_code}) break except requests.exceptions.RequestException as e: print(f网络错误{e}) retry_count 1 time.sleep(2) if response is not None and response.status_code 200: result response.json() print(result)这里的timeout30是请求级超时控制防止某个批次把主流程拖死。重试侧面分别处理超时和 HTTP 错误尤其是 429限流和 5xx服务端异常这两类故障重试才有效果4xx 多半是参数问题重试只是浪费时间。2.3 响应解析与批次结果的顺序保持DeepSeek 标准聊天补全接口的返回结构里结果落在choices[0].message.content。如果每个批量元素都有独立id我会在解析时把原始 id 和模型输出拼成一个可追踪的结构if response.status_code 200: parsed [] resp_json response.json() # 单条聊天补全场景模型返回在 choices[0].message.content content resp_json[choices][0][message][content] parsed.append({id: 1, result: content})这里要关注两个问题一是结果顺序是否与入参一致绝大多数接口是同步返回所以顺序一致但稳妥做法是给每个子请求带上id解析时按id对应回去不要依赖数组下标。二是响应体里的usage字段包含 prompt_tokens 和 completion_tokens我会在批量场景里把它累加统计用来估算成本也能反过来判断某批结果是否异常短小。2.4 批量请求的边界什么情况不适合硬合并批量不是越大越好。单批 payload 超出服务端大小限制时会收到 413单批里混杂了完全不同的任务类型比如一批做摘要、一批做分类提示词不同会导致无法统一复用同一个 system prompt需要按任务类型拆成多个批次。我在实际项目里的习惯是把任务先按“模型参数一致 系统提示一致”分组每组内部再按 10 到 50 条拆批而不是拿到数据直接一把梭。3. 异步调用用事件循环把等待时间抢回来3.1 为什么 I/O 密集场景下 Async 比线程池更划算同步请求的最大问题是阻塞你发一个请求后程序就干等着网络返回CPU 明明空闲却什么也做不了。线程池能解决一部分问题但线程切换有成本线程数量一多调度开销反而吃掉收益。asyncio的做法是在单线程里用事件循环管理大量协程遇到await就把控制权交出去让程序在等待网络响应时去处理其他请求。DeepSeek 调用是典型的 I/O 密集型任务瓶颈在网络往返和模型推理等待CPU 本地计算非常少。这个问题上用 asyncio不需要管线程锁也不用担心线程池炸掉内存一套事件循环就能撑起几十上百个并发请求。Python 3.7 之后用asyncio.run()启动入口配合aiohttp做 HTTP 客户端是当前最顺手的组合。3.2 最小可运行的异步调用 DeepSeek 实现先装aiohttppip install aiohttp然后实现一个异步调用函数import asyncio import aiohttp async def call_deepseek_async(session, payload): headers { Content-Type: application/json, Authorization: fBearer {API_KEY} } try: async with session.post(API_URL, headersheaders, jsonpayload, timeout30) as resp: if resp.status 200: return await resp.json() else: print(fHTTP {resp.status}) return None except Exception as e: print(f调用异常{e}) return None async def main(): batch_data [ {model: deepseek-chat, messages: [...], temperature: 0.2} for _ in range(10) ] async with aiohttp.ClientSession() as session: tasks [call_deepseek_async(session, p) for p in batch_data] results await asyncio.gather(*tasks) valid [r for r in results if r is not None] print(f成功 {len(valid)} / {len(batch_data)})asyncio.gather会并发执行所有传入的协程并在所有任务完成后返回结果列表。这里有个容易忽略的细节aiohttp.ClientSession需要复用同一个实例不能每个请求创建一次 session否则底层连接池形同虚设并发收益会打折。3.3 批量 异步结合顺序保存与异常隔离把前面两章的内容拼在一起时最常见的错误是只同步返回结果数组不检查单个协程的异常状态。gather默认遇到第一个异常就会向上抛出导致已经成功的请求结果全部丢失。我一般会加上return_exceptionsTrueresults await asyncio.gather(*tasks, return_exceptionsTrue)这样每个失败的协程会返回一个异常对象而不是直接中断整个流程。随后遍历results对isinstance(r, Exception)的结果走单独的重试或落盘记录成功的再解析内容。顺序问题同样靠入参里的id字段兜底不要让结果处理逻辑假设数组顺序等于入参顺序。4. 并发控制与重试吞吐量能不能稳住就看这两个参数4.1 用 Semaphore 限制并发数并发数不是越大越好。开 200 个协程同时打账户限流一触发全是 429重试又叠加新的并发反而把服务端打挂。我一般第一版先限制并发 3 到 5跑通后再逐步往上加。sem asyncio.Semaphore(5) async def limited_call(session, payload): async with sem: return await call_deepseek_async(session, payload)Semaphore(5)表示同一时刻最多 5 个协程进入async with sem内部其他协程在门口排队。这个参数本质上是给账户限流留缓冲同时避免本地文件句柄和内存被占满。实际取值没有标准答案我通常先看单次请求的平均耗时如果平均要 1 秒并发 10 大约能撑每秒 10 次请求如果账户配额是每分钟 60 次那并发 5 已经足够。4.2 指数退避重试429 和 5xx 分开处理重试策略最怕的是“失败就立刻重试”这在高并发下会形成重试风暴。经验做法是429 用指数退避加随机抖动5xx 可以适当重试4xx 一律不重试。import random async def call_with_retry(session, payload, max_retries3): for attempt in range(max_retries): try: async with session.post(API_URL, headersheaders, jsonpayload, timeout30) as resp: if resp.status 200: return await resp.json() if resp.status in (429, 500, 502, 503): delay 2 ** attempt random.uniform(0, 1) print(fHTTP {resp.status}等待 {delay:.2f}s) await asyncio.sleep(delay) else: return None except Exception as e: delay 2 ** attempt random.uniform(0, 1) await asyncio.sleep(delay) return None2 ** attempt是重试基础间隔第 0 次重试等 1 秒第 1 次等 2 秒第二次等 4 秒。random.uniform(0, 1)加抖动是为了让多个协程的重试时间错开避免所有请求同时醒来重新打向服务器。这个方法在线上效果很直接我几次翻车都是没加抖动大量重试请求在同一秒涌进去导致限流时间被拉长。4.3 动态调整并发小数据探底、中数据验证、全量放行固定并发数能解决大部分问题但遇到一天要跑几十万条数据时固定值容易偏保守。我一般按三步走先用 50 条数据、并发 5 跑一轮看错误率和平均耗时再把并发调到 10、20 各跑一轮记录错误率拐点最后在错误率低于 1% 的前提下用拐点并发跑全量。所谓“动态调整”不需要上什么自动化框架把并发数做成一个命令行参数每次跑前手动验证一次就行比写一堆自适应逻辑更可控。5. 避坑与排错DeepSeek 批量异步调用中的五个典型问题5.1 请求超时重试后整体耗时反而翻倍现象单批请求 30 秒超时重试 3 次后成功但整批任务耗时从 3 分钟涨到 15 分钟。原因并发设置过高单次请求排队时间已经超过客户端超时阈值。服务端没拒绝请求只是处理不过来重试让同一批数据反复排队。解决把并发数降下来同时把超时从 30 秒提到 60 秒。优先保证请求不超时再用重试兜底。我踩过这个坑之后都把超时和并发当成一对参数调并发翻了倍超时必须同步放大。5.2 批量请求返回 413payload 太大被服务端拒绝现象一批塞了 100 条文本请求直接返回 413 Request Entity Too Large。原因单批请求体超过服务端限制常见于长文本场景。解决按字符数或条数拆批。我一般先估算单条请求体平均大小然后控制单批总大小在 1MB 以内长文本场景每批 10 条起步短文本可以到 50 条。拆分逻辑写成函数传入文本列表和批次上限自动切分后再逐批处理。5.3 异步任务失败但日志里什么都没有现象gather 返回后结果缺失也没有任何错误输出。原因协程内部把异常吞掉了。早期版本里我在 except 块里只打印e但打印本身放在return None之前异常对象被覆盖后续排查拿不到任何有效信息。解决异常处理里先记录完整堆栈再决定返回默认值还是重试。使用traceback.format_exc()把堆栈存到日志不要把 exception 对象直接转字符串。从那以后我每个异步调用函数都强制要求要么抛出可重试异常要么返回带错误码的结果对象不允许静默返回 None。5.4 本地 Windows 环境跑 async 代码报事件循环错误现象代码在 Linux 上正常Windows 上运行时报RuntimeError: Event loop is closed。原因Windows 下asyncio的事件循环策略和 Proactor 模式有差异某些 Python 版本在重复调用asyncio.run()时会触发旧事件循环残留问题。解决入口只调用一次asyncio.run(main())不要在循环里反复创建事件循环。如果代码被外部框架调用使用asyncio.get_event_loop()配合run_until_complete()兜底而不是在函数内部擅自创建新事件循环。5.5 鉴权失败 401Authorization 头拼接方式不对现象批量请求全部返回 401检查密钥没发现问题。原因Authorization 头必须是Bearer 密钥的形式中间有空格或者密钥被去掉了前缀。另外如果在代理环境下请求头信息可能被代理改写。解决先打印请求头的实际值确认格式再确认代理环境不会改写鉴权头。排查时我会用一条最小请求单独测省得混在批量里难定位。6. 用一套压测脚本验证“10 倍”同步循环与异步批量的对比基准6.1 快速对比压测脚本要验证是否真有 10 倍提升不要凭感觉直接跑一轮同数据量的对比。压测脚本的核心逻辑是同样的任务集合分别用同步 for 循环和异步 Semaphore 调度执行记录总耗时、成功率和错误分布。import asyncio import time def run_sync(tasks_data): start time.time() results [] for item in tasks_data: resp requests.post(API_URL, headersheaders, jsonitem, timeout30) results.append(resp.status_code) return time.time() - start, results async def run_async(tasks_data, concurrency5): start time.time() sem asyncio.Semaphore(concurrency) async with aiohttp.ClientSession() as session: async def one(item): async with sem: return await call_with_retry(session, item) results await asyncio.gather(*[one(x) for x in tasks_data]) return time.time() - start, results # 用同一批任务对比 sync_time, sync_res run_sync(batch_data) async_time, async_res asyncio.run(run_async(batch_data)) print(fsync{sync_time:.2f}s, async{async_time:.2f}s, 提升{sync_time/async_time:.1f}x)这里每一步都只统计了耗时和状态码分布。真正压测时我还会记录错误码明细单独统计 429 重试次数。如果异步跑下来提升不到 3 倍先看并发是不是设太低再看是不是本地网络或代理成了新瓶颈。6.2 怎么读压测结果压测结果里的平均耗时会被个别慢请求抬高我更关注 P95 和错误率。P95 越小说明大部分请求的响应时间稳定错误率如果超过 1%优先降并发而不是加超时。另外要对比的是重试次数重试多说明并发已经逼近限流阈值继续提高并发只会让调整个更慢。6.3 我从那次之后固定的三步走现在每次接 DeepSeek 相关任务我都强制先过一遍小样本验证再放量。第一步用 50 条数据配并发 5 跑通流程确认鉴权、解析、落库全链路正常第二步用 500 条数据调并发 10 和 20记录错误率拐点第三步才用全量数据开盘跑。开盘时保留并发数参数线上运行中如果错误率突然上升先降一半并发而不是立即调大概率重试等待。这套流程把“能不能跑”和“能跑多快”分开验证数据清洗、评论分类、批量写摘要这些场景都能套用。那次 40 分钟压到 3 分半的改造之后我再也没用裸 for 循环去调任何大模型 API凡是超过 100 条的任务第一反应就是先拆批、再上并发。希望这套拆解和踩坑记录帮得到你。本文还有配套的精品资源点击获取