Qwen2.5-72B-Instruct-GPTQ-Int4实战:批量API调用与并发压力测试

1. 引言:从单次对话到批量处理

如果你已经成功部署了Qwen2.5-72B-Instruct-GPTQ-Int4,并且通过Chainlit前端和它愉快地聊过天,那么恭喜你,你已经迈出了第一步。但你可能也发现了,每次只能问一个问题,等它回答完才能问下一个。这就像你有一台性能强大的服务器,却只用来浏览网页,是不是有点大材小用?

在实际的生产环境或研究场景中,我们往往需要处理大量的文本生成任务。比如:

  • 批量处理成百上千条客服对话,生成标准回复
  • 对海量文档进行摘要提取或关键词生成
  • 同时为多个用户提供问答服务
  • 测试模型在高并发下的稳定性和响应速度

这时候,我们就需要从“单兵作战”切换到“集团军作战”模式。今天这篇文章,我就带你深入实战,看看如何对Qwen2.5-72B-Instruct-GPTQ-Int4进行批量API调用,并进行并发压力测试,真正把这头“巨兽”的性能压榨出来。

2. 环境准备与API接口确认

在开始批量调用之前,我们首先要确认模型服务已经正常运行,并且知道如何通过API与它通信。

2.1 确认模型服务状态

还记得部署时查看日志的方法吗?我们再来确认一下服务是否正常:

# 查看模型服务日志
cat /root/workspace/llm.log

如果看到类似下面的输出,说明模型已经加载成功并正在运行:

INFO:     Started server process [12345]
INFO:     Waiting for application startup.
INFO:     Application startup complete.
INFO:     Uvicorn running on http://0.0.0.0:8000 (Press CTRL+C to quit)

关键信息是Uvicorn running on http://0.0.0.0:8000,这告诉我们API服务运行在8000端口。

2.2 了解vLLM的API接口

vLLM提供了标准的OpenAI兼容的API接口,这意味着我们可以用和调用ChatGPT API几乎相同的方式来调用我们的本地模型。主要接口有两个:

  1. 聊天补全接口 (/v1/chat/completions) - 最常用的接口,支持对话格式
  2. 补全接口 (/v1/completions) - 简单的文本补全

对于Qwen2.5-72B-Instruct这种指令调优模型,我们主要使用聊天补全接口。它的请求格式是这样的:

{
    "model": "Qwen2.5-72B-Instruct-GPTQ-Int4",
    "messages": [
        {"role": "system", "content": "你是一个有帮助的助手"},
        {"role": "user", "content": "你好,请介绍一下你自己"}
    ],
    "temperature": 0.7,
    "max_tokens": 512
}

3. 单次API调用:从Chainlit到Python代码

在Chainlit里点点按钮就能聊天,背后其实也是通过API调用的。现在让我们看看如何用Python代码实现同样的功能。

3.1 安装必要的Python库

首先确保你安装了requests库,如果没有的话:

pip install requests

3.2 编写第一个API调用脚本

创建一个简单的Python脚本,测试单次API调用:

import requests
import json
import time

def single_api_call(prompt, system_prompt="你是一个有帮助的助手"):
    """
    单次API调用函数
    """
    # API端点
    url = "http://localhost:8000/v1/chat/completions"
    
    # 请求头
    headers = {
        "Content-Type": "application/json"
    }
    
    # 请求体
    payload = {
        "model": "Qwen2.5-72B-Instruct-GPTQ-Int4",
        "messages": [
            {"role": "system", "content": system_prompt},
            {"role": "user", "content": prompt}
        ],
        "temperature": 0.7,
        "max_tokens": 512,
        "stream": False  # 非流式响应,一次性返回完整结果
    }
    
    try:
        # 记录开始时间
        start_time = time.time()
        
        # 发送请求
        response = requests.post(url, headers=headers, data=json.dumps(payload))
        
        # 记录结束时间
        end_time = time.time()
        
        if response.status_code == 200:
            result = response.json()
            content = result["choices"][0]["message"]["content"]
            tokens_used = result["usage"]["total_tokens"]
            response_time = end_time - start_time
            
            print(f"✅ 请求成功!")
            print(f"📝 问题:{prompt}")
            print(f"💬 回答:{content}")
            print(f"🔢 使用token数:{tokens_used}")
            print(f"⏱️  响应时间:{response_time:.2f}秒")
            print("-" * 50)
            
            return {
                "success": True,
                "content": content,
                "tokens": tokens_used,
                "response_time": response_time
            }
        else:
            print(f"❌ 请求失败,状态码:{response.status_code}")
            print(f"错误信息:{response.text}")
            return {"success": False, "error": response.text}
            
    except Exception as e:
        print(f"❌ 请求异常:{str(e)}")
        return {"success": False, "error": str(e)}

# 测试单次调用
if __name__ == "__main__":
    test_prompt = "用简单的语言解释什么是人工智能"
    result = single_api_call(test_prompt)

运行这个脚本,你应该能看到类似Chainlit界面的回答,同时还能看到响应时间和token使用量。这就是API调用的基础。

4. 批量API调用:处理多个任务

单次调用没问题了,现在我们来处理批量任务。批量调用有两种常见方式:顺序调用和并发调用。

4.1 顺序批量调用

顺序调用就是一个个来,等上一个完成再处理下一个。这种方法简单直接,适合任务量不大或者对并发要求不高的场景。

import requests
import json
import time
from typing import List, Dict

def sequential_batch_calls(prompts: List[str], system_prompt="你是一个有帮助的助手"):
    """
    顺序批量调用API
    """
    url = "http://localhost:8000/v1/chat/completions"
    headers = {"Content-Type": "application/json"}
    
    results = []
    total_tokens = 0
    total_time = 0
    
    print(f"🚀 开始顺序处理 {len(prompts)} 个任务...")
    print("=" * 60)
    
    for i, prompt in enumerate(prompts, 1):
        print(f"📋 处理第 {i}/{len(prompts)} 个任务...")
        
        payload = {
            "model": "Qwen2.5-72B-Instruct-GPTQ-Int4",
            "messages": [
                {"role": "system", "content": system_prompt},
                {"role": "user", "content": prompt}
            ],
            "temperature": 0.7,
            "max_tokens": 256,
            "stream": False
        }
        
        try:
            start_time = time.time()
            response = requests.post(url, headers=headers, data=json.dumps(payload))
            end_time = time.time()
            
            if response.status_code == 200:
                result = response.json()
                content = result["choices"][0]["message"]["content"]
                tokens = result["usage"]["total_tokens"]
                response_time = end_time - start_time
                
                total_tokens += tokens
                total_time += response_time
                
                results.append({
                    "prompt": prompt,
                    "content": content,
                    "tokens": tokens,
                    "response_time": response_time,
                    "success": True
                })
                
                print(f"   ✅ 成功 | 用时:{response_time:.2f}s | Token:{tokens}")
            else:
                results.append({
                    "prompt": prompt,
                    "content": None,
                    "tokens": 0,
                    "response_time": 0,
                    "success": False,
                    "error": f"HTTP {response.status_code}"
                })
                print(f"   ❌ 失败 | 错误:HTTP {response.status_code}")
                
        except Exception as e:
            results.append({
                "prompt": prompt,
                "content": None,
                "tokens": 0,
                "response_time": 0,
                "success": False,
                "error": str(e)
            })
            print(f"   ❌ 异常 | 错误:{str(e)}")
        
        # 为了避免请求过于密集,可以添加小延迟
        time.sleep(0.1)
    
    print("=" * 60)
    print(f"🎯 任务完成统计:")
    print(f"   成功:{sum(1 for r in results if r['success'])}/{len(prompts)}")
    print(f"   总用时:{total_time:.2f}秒")
    print(f"   平均响应时间:{total_time/len(prompts):.2f}秒")
    print(f"   总Token使用:{total_tokens}")
    print(f"   平均Token/请求:{total_tokens/len(prompts):.0f}")
    
    return results

# 测试顺序批量调用
if __name__ == "__main__":
    # 准备测试问题
    test_prompts = [
        "用一句话解释机器学习",
        "Python和JavaScript的主要区别是什么?",
        "如何煮一碗好吃的泡面?",
        "地球到月球的平均距离是多少?",
        "推荐三本值得读的科幻小说"
    ]
    
    results = sequential_batch_calls(test_prompts)
    
    # 保存结果到文件
    with open("sequential_results.json", "w", encoding="utf-8") as f:
        json.dump(results, f, ensure_ascii=False, indent=2)
    print("💾 结果已保存到 sequential_results.json")

顺序调用的优点是实现简单,不会给服务器造成太大压力。缺点是速度慢,总时间等于所有请求时间的总和。

4.2 并发批量调用

当我们需要处理大量任务时,顺序调用就太慢了。这时候就需要并发调用,同时发送多个请求。Python中我们可以用concurrent.futures模块来实现。

import requests
import json
import time
from concurrent.futures import ThreadPoolExecutor, as_completed
from typing import List, Dict

def make_api_request(prompt: str, system_prompt: str = "你是一个有帮助的助手"):
    """
    单个API请求函数,用于并发调用
    """
    url = "http://localhost:8000/v1/chat/completions"
    headers = {"Content-Type": "application/json"}
    
    payload = {
        "model": "Qwen2.5-72B-Instruct-GPTQ-Int4",
        "messages": [
            {"role": "system", "content": system_prompt},
            {"role": "user", "content": prompt}
        ],
        "temperature": 0.7,
        "max_tokens": 256,
        "stream": False
    }
    
    try:
        start_time = time.time()
        response = requests.post(url, headers=headers, data=json.dumps(payload))
        end_time = time.time()
        
        if response.status_code == 200:
            result = response.json()
            return {
                "prompt": prompt,
                "content": result["choices"][0]["message"]["content"],
                "tokens": result["usage"]["total_tokens"],
                "response_time": end_time - start_time,
                "success": True
            }
        else:
            return {
                "prompt": prompt,
                "content": None,
                "tokens": 0,
                "response_time": 0,
                "success": False,
                "error": f"HTTP {response.status_code}"
            }
    except Exception as e:
        return {
            "prompt": prompt,
            "content": None,
            "tokens": 0,
            "response_time": 0,
            "success": False,
            "error": str(e)
        }

def concurrent_batch_calls(prompts: List[str], max_workers: int = 5, system_prompt="你是一个有帮助的助手"):
    """
    并发批量调用API
    max_workers: 最大并发数,根据你的服务器配置调整
    """
    print(f"🚀 开始并发处理 {len(prompts)} 个任务,最大并发数:{max_workers}")
    print("=" * 60)
    
    results = []
    start_total_time = time.time()
    
    # 使用线程池并发执行
    with ThreadPoolExecutor(max_workers=max_workers) as executor:
        # 提交所有任务
        future_to_prompt = {
            executor.submit(make_api_request, prompt, system_prompt): prompt 
            for prompt in prompts
        }
        
        # 处理完成的任务
        for i, future in enumerate(as_completed(future_to_prompt), 1):
            prompt = future_to_prompt[future]
            try:
                result = future.result()
                results.append(result)
                
                if result["success"]:
                    print(f"📋 {i}/{len(prompts)} ✅ 成功 | 用时:{result['response_time']:.2f}s | Token:{result['tokens']}")
                else:
                    print(f"📋 {i}/{len(prompts)} ❌ 失败 | 错误:{result.get('error', '未知错误')}")
            except Exception as e:
                print(f"📋 {i}/{len(prompts)} ❌ 异常 | 错误:{str(e)}")
                results.append({
                    "prompt": prompt,
                    "content": None,
                    "tokens": 0,
                    "response_time": 0,
                    "success": False,
                    "error": str(e)
                })
    
    end_total_time = time.time()
    total_duration = end_total_time - start_total_time
    
    # 统计结果
    successful = [r for r in results if r["success"]]
    failed = [r for r in results if not r["success"]]
    
    total_tokens = sum(r["tokens"] for r in successful)
    total_response_time = sum(r["response_time"] for r in successful)
    
    print("=" * 60)
    print(f"🎯 并发任务完成统计:")
    print(f"   成功:{len(successful)}/{len(prompts)}")
    print(f"   失败:{len(failed)}/{len(prompts)}")
    print(f"   总耗时:{total_duration:.2f}秒")
    print(f"   总响应时间(所有请求之和):{total_response_time:.2f}秒")
    
    if successful:
        avg_response_time = total_response_time / len(successful)
        avg_tokens = total_tokens / len(successful)
        print(f"   平均响应时间:{avg_response_time:.2f}秒")
        print(f"   总Token使用:{total_tokens}")
        print(f"   平均Token/请求:{avg_tokens:.0f}")
        
        # 计算并发效率
        if len(successful) > 0:
            efficiency = total_response_time / total_duration if total_duration > 0 else 0
            print(f"   并发效率:{efficiency:.2f}(越接近并发数越好)")
    
    return results

# 测试并发批量调用
if __name__ == "__main__":
    # 准备更多测试问题
    test_prompts = [
        "解释什么是深度学习",
        "如何学习编程?给一些建议",
        "中国的首都是哪里?",
        "写一个简单的Python函数计算斐波那契数列",
        "什么是气候变化?",
        "推荐几个适合初学者的编程项目",
        "如何保持健康的生活方式?",
        "解释区块链技术的基本原理",
        "什么是量子计算?",
        "如何提高英语口语能力?"
    ]
    
    # 测试不同并发数
    for workers in [3, 5, 8]:
        print(f"\n🔧 测试并发数:{workers}")
        print("-" * 40)
        results = concurrent_batch_calls(test_prompts, max_workers=workers)
        
        # 保存结果
        with open(f"concurrent_{workers}workers_results.json", "w", encoding="utf-8") as f:
            json.dump(results, f, ensure_ascii=False, indent=2)
        print(f"💾 结果已保存到 concurrent_{workers}workers_results.json")

并发调用的速度比顺序调用快得多,但要注意控制并发数,避免把服务器压垮。

5. 压力测试:探索性能极限

现在我们来做个真正的压力测试,看看Qwen2.5-72B-Instruct-GPTQ-Int4在高并发下的表现。

5.1 设计压力测试方案

一个好的压力测试应该考虑以下几个方面:

  1. 并发数梯度:从低到高测试不同并发数
  2. 请求频率:模拟不同的请求间隔
  3. 请求内容:使用不同长度和复杂度的提示词
  4. 持续时间:测试长时间运行稳定性
  5. 监控指标:响应时间、成功率、错误率、Token使用量

5.2 实现压力测试脚本

import requests
import json
import time
import random
from concurrent.futures import ThreadPoolExecutor, as_completed
from threading import Lock
from datetime import datetime
import statistics

class PressureTester:
    def __init__(self, base_url="http://localhost:8000"):
        self.base_url = base_url
        self.results = []
        self.lock = Lock()  # 用于线程安全地更新结果
        self.test_start_time = None
        self.test_end_time = None
        
    def generate_test_prompts(self, count=50):
        """
        生成测试用的提示词
        包含不同长度和复杂度的提示
        """
        prompts = []
        
        # 短提示
        short_prompts = [
            "你好",
            "今天天气怎么样?",
            "现在几点了?",
            "谢谢",
            "再见"
        ]
        
        # 中等长度提示
        medium_prompts = [
            "用简单的语言解释什么是机器学习",
            "Python和Java的主要区别是什么?",
            "如何煮一碗好吃的面条?",
            "地球到太阳的距离是多少?",
            "推荐三本值得读的书籍"
        ]
        
        # 长提示
        long_prompts = [
            "请详细解释深度学习中的卷积神经网络是如何工作的,包括卷积层、池化层和全连接层的作用,以及为什么它在图像识别任务中表现优异。",
            "写一篇关于人工智能对社会影响的短文,讨论其积极和消极方面,以及我们应该如何应对这些挑战。字数在300字左右。",
            "请为一家新开的咖啡店设计一个营销方案,包括目标客户分析、产品定位、促销活动和社交媒体策略。要求具体可行。"
        ]
        
        # 混合生成测试提示
        for i in range(count):
            # 随机选择提示类型
            prompt_type = random.choice(["short", "medium", "long"])
            
            if prompt_type == "short":
                prompt = random.choice(short_prompts)
            elif prompt_type == "medium":
                prompt = random.choice(medium_prompts)
            else:
                prompt = random.choice(long_prompts)
            
            # 添加序号以便追踪
            prompts.append(f"[{i+1}] {prompt}")
        
        return prompts
    
    def make_request(self, prompt_id, prompt, system_prompt="你是一个有帮助的助手"):
        """
        发送单个API请求
        """
        url = f"{self.base_url}/v1/chat/completions"
        headers = {"Content-Type": "application/json"}
        
        # 随机化一些参数,模拟真实场景
        temperature = random.uniform(0.5, 0.9)
        max_tokens = random.choice([128, 256, 512])
        
        payload = {
            "model": "Qwen2.5-72B-Instruct-GPTQ-Int4",
            "messages": [
                {"role": "system", "content": system_prompt},
                {"role": "user", "content": prompt}
            ],
            "temperature": temperature,
            "max_tokens": max_tokens,
            "stream": False
        }
        
        result = {
            "prompt_id": prompt_id,
            "prompt": prompt,
            "temperature": temperature,
            "max_tokens": max_tokens,
            "success": False,
            "response_time": 0,
            "tokens": 0,
            "error": None,
            "timestamp": time.time()
        }
        
        try:
            start_time = time.time()
            response = requests.post(url, headers=headers, data=json.dumps(payload), timeout=60)
            end_time = time.time()
            
            result["response_time"] = end_time - start_time
            
            if response.status_code == 200:
                response_data = response.json()
                result["success"] = True
                result["tokens"] = response_data["usage"]["total_tokens"]
                result["content_length"] = len(response_data["choices"][0]["message"]["content"])
            else:
                result["success"] = False
                result["error"] = f"HTTP {response.status_code}: {response.text[:100]}"
                
        except requests.exceptions.Timeout:
            result["success"] = False
            result["error"] = "请求超时(60秒)"
        except requests.exceptions.ConnectionError:
            result["success"] = False
            result["error"] = "连接错误"
        except Exception as e:
            result["success"] = False
            result["error"] = str(e)
        
        # 线程安全地保存结果
        with self.lock:
            self.results.append(result)
        
        return result
    
    def run_pressure_test(self, total_requests=100, concurrent_workers=10, request_interval=0.1):
        """
        运行压力测试
        total_requests: 总请求数
        concurrent_workers: 并发工作线程数
        request_interval: 请求间隔(秒),0表示无间隔
        """
        print("=" * 70)
        print("🔥 Qwen2.5-72B-Instruct-GPTQ-Int4 压力测试开始")
        print("=" * 70)
        print(f"测试配置:")
        print(f"  总请求数:{total_requests}")
        print(f"  并发数:{concurrent_workers}")
        print(f"  请求间隔:{request_interval}秒")
        print(f"  开始时间:{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}")
        print("-" * 70)
        
        # 生成测试提示
        prompts = self.generate_test_prompts(total_requests)
        
        # 重置结果
        self.results = []
        self.test_start_time = time.time()
        
        completed_requests = 0
        successful_requests = 0
        failed_requests = 0
        
        # 使用线程池进行并发测试
        with ThreadPoolExecutor(max_workers=concurrent_workers) as executor:
            # 提交所有任务
            future_to_id = {}
            for i, prompt in enumerate(prompts):
                future = executor.submit(self.make_request, i+1, prompt)
                future_to_id[future] = i+1
                
                # 控制请求发送频率
                if request_interval > 0:
                    time.sleep(request_interval)
            
            # 处理完成的任务
            for future in as_completed(future_to_id):
                prompt_id = future_to_id[future]
                try:
                    result = future.result()
                    completed_requests += 1
                    
                    if result["success"]:
                        successful_requests += 1
                        status = "✅"
                    else:
                        failed_requests += 1
                        status = "❌"
                    
                    # 实时显示进度
                    progress = completed_requests / total_requests * 100
                    print(f"[{status}] 请求 {prompt_id}/{total_requests} ({progress:.1f}%) | "
                          f"用时:{result['response_time']:.2f}s | "
                          f"Token:{result['tokens'] if result['success'] else 'N/A'} | "
                          f"错误:{result['error'] or '无'}")
                    
                except Exception as e:
                    print(f"[❌] 请求 {prompt_id} 异常:{str(e)}")
                    completed_requests += 1
                    failed_requests += 1
        
        self.test_end_time = time.time()
        total_duration = self.test_end_time - self.test_start_time
        
        # 生成测试报告
        self.generate_report(total_requests, successful_requests, failed_requests, total_duration)
        
        return self.results
    
    def generate_report(self, total_requests, successful_requests, failed_requests, total_duration):
        """
        生成压力测试报告
        """
        print("\n" + "=" * 70)
        print("📊 压力测试报告")
        print("=" * 70)
        
        # 基础统计
        success_rate = successful_requests / total_requests * 100 if total_requests > 0 else 0
        
        print(f"📈 基础统计:")
        print(f"  总请求数:{total_requests}")
        print(f"  成功请求:{successful_requests}")
        print(f"  失败请求:{failed_requests}")
        print(f"  成功率:{success_rate:.1f}%")
        print(f"  总测试时长:{total_duration:.2f}秒")
        print(f"  平均QPS(每秒查询数):{total_requests/total_duration:.2f}")
        
        # 响应时间统计
        successful_results = [r for r in self.results if r["success"]]
        if successful_results:
            response_times = [r["response_time"] for r in successful_results]
            token_counts = [r["tokens"] for r in successful_results]
            
            print(f"\n⏱️  响应时间分析:")
            print(f"  平均响应时间:{statistics.mean(response_times):.2f}秒")
            print(f"  中位数响应时间:{statistics.median(response_times):.2f}秒")
            print(f"  最快响应时间:{min(response_times):.2f}秒")
            print(f"  最慢响应时间:{max(response_times):.2f}秒")
            print(f"  95%请求响应时间:{sorted(response_times)[int(len(response_times)*0.95)]:.2f}秒")
            
            # 响应时间分布
            time_buckets = {"<1s": 0, "1-3s": 0, "3-5s": 0, "5-10s": 0, ">10s": 0}
            for rt in response_times:
                if rt < 1:
                    time_buckets["<1s"] += 1
                elif rt < 3:
                    time_buckets["1-3s"] += 1
                elif rt < 5:
                    time_buckets["3-5s"] += 1
                elif rt < 10:
                    time_buckets["5-10s"] += 1
                else:
                    time_buckets[">10s"] += 1
            
            print(f"\n📊 响应时间分布:")
            for bucket, count in time_buckets.items():
                percentage = count / len(response_times) * 100
                print(f"  {bucket}: {count}次 ({percentage:.1f}%)")
            
            print(f"\n🔢 Token使用统计:")
            print(f"  总Token使用:{sum(token_counts)}")
            print(f"  平均Token/请求:{statistics.mean(token_counts):.0f}")
            print(f"  最少Token:{min(token_counts)}")
            print(f"  最多Token:{max(token_counts)}")
        
        # 错误分析
        failed_results = [r for r in self.results if not r["success"]]
        if failed_results:
            print(f"\n❌ 错误分析:")
            error_types = {}
            for r in failed_results:
                error = r["error"]
                error_types[error] = error_types.get(error, 0) + 1
            
            for error, count in error_types.items():
                print(f"  {error}: {count}次")
        
        print(f"\n⏰ 测试结束时间:{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}")
        print("=" * 70)
        
        # 保存详细结果
        timestamp = datetime.now().strftime("%Y%m%d_%H%M%S")
        filename = f"pressure_test_results_{timestamp}.json"
        
        report_data = {
            "test_config": {
                "total_requests": total_requests,
                "concurrent_workers": getattr(self, 'current_workers', 'N/A'),
                "request_interval": getattr(self, 'current_interval', 'N/A'),
                "start_time": datetime.fromtimestamp(self.test_start_time).strftime('%Y-%m-%d %H:%M:%S'),
                "end_time": datetime.fromtimestamp(self.test_end_time).strftime('%Y-%m-%d %H:%M:%S'),
                "total_duration": total_duration
            },
            "summary": {
                "total_requests": total_requests,
                "successful_requests": successful_requests,
                "failed_requests": failed_requests,
                "success_rate": success_rate,
                "avg_qps": total_requests/total_duration if total_duration > 0 else 0
            },
            "response_time_stats": {
                "mean": statistics.mean(response_times) if successful_results else 0,
                "median": statistics.median(response_times) if successful_results else 0,
                "min": min(response_times) if successful_results else 0,
                "max": max(response_times) if successful_results else 0,
                "p95": sorted(response_times)[int(len(response_times)*0.95)] if successful_results else 0
            } if successful_results else {},
            "token_stats": {
                "total": sum(token_counts) if successful_results else 0,
                "mean": statistics.mean(token_counts) if successful_results else 0,
                "min": min(token_counts) if successful_results else 0,
                "max": max(token_counts) if successful_results else 0
            } if successful_results else {},
            "detailed_results": self.results
        }
        
        with open(filename, "w", encoding="utf-8") as f:
            json.dump(report_data, f, ensure_ascii=False, indent=2)
        
        print(f"💾 详细结果已保存到:{filename}")

# 运行压力测试
if __name__ == "__main__":
    tester = PressureTester()
    
    # 测试不同配置
    test_configs = [
        {"total_requests": 50, "concurrent_workers": 5, "request_interval": 0.2},
        {"total_requests": 100, "concurrent_workers": 10, "request_interval": 0.1},
        {"total_requests": 50, "concurrent_workers": 20, "request_interval": 0.05},
    ]
    
    for config in test_configs:
        print(f"\n{'='*70}")
        print(f"🔧 测试配置:{config['total_requests']}请求,{config['concurrent_workers']}并发,间隔{config['request_interval']}秒")
        print(f"{'='*70}")
        
        # 设置当前配置(用于报告)
        tester.current_workers = config["concurrent_workers"]
        tester.current_interval = config["request_interval"]
        
        # 运行测试
        results = tester.run_pressure_test(
            total_requests=config["total_requests"],
            concurrent_workers=config["concurrent_workers"],
            request_interval=config["request_interval"]
        )
        
        # 每次测试后休息一下,让服务器恢复
        print("😴 测试完成,等待10秒让服务器恢复...")
        time.sleep(10)

这个压力测试脚本可以帮你全面了解模型的性能表现。你可以根据实际情况调整测试参数。

6. 性能优化与最佳实践

通过压力测试,你可能会发现一些性能瓶颈。这里分享一些优化建议:

6.1 调整vLLM配置

vLLM有很多配置参数可以优化性能。如果你有权限修改部署配置,可以尝试:

# vLLM部署时的优化参数示例
# 在启动vLLM时添加这些参数

# 增加批处理大小,提高吞吐量
--max_num_batched_tokens 4096

# 调整KV缓存大小
--block_size 16

# 启用流水线并行(如果有多GPU)
--tensor-parallel-size 2
--pipeline-parallel-size 2

# 使用PagedAttention优化内存
--enable-paged-attention

6.2 客户端优化建议

  1. 连接池管理:重用HTTP连接,避免频繁建立连接的开销
  2. 请求批处理:将多个请求合并为一个批次发送
  3. 超时设置:合理设置请求超时时间
  4. 错误重试:实现指数退避的重试机制
  5. 限流控制:根据服务器能力控制请求频率

6.3 监控与告警

在生产环境中,建议添加监控:

import psutil
import time
from datetime import datetime

def monitor_system_resources(interval=5, duration=300):
    """
    监控系统资源使用情况
    """
    print("🖥️  开始监控系统资源...")
    print("时间 | CPU使用率 | 内存使用率 | 已用内存(GB)")
    print("-" * 50)
    
    start_time = time.time()
    while time.time() - start_time < duration:
        # CPU使用率
        cpu_percent = psutil.cpu_percent(interval=1)
        
        # 内存使用情况
        memory = psutil.virtual_memory()
        memory_percent = memory.percent
        memory_used_gb = memory.used / (1024**3)
        
        # 当前时间
        current_time = datetime.now().strftime("%H:%M:%S")
        
        print(f"{current_time} | {cpu_percent:6.1f}% | {memory_percent:6.1f}% | {memory_used_gb:8.2f} GB")
        
        time.sleep(interval)
    
    print("📊 监控结束")

# 在压力测试时同时监控系统资源
if __name__ == "__main__":
    import threading
    
    # 启动监控线程
    monitor_thread = threading.Thread(target=monitor_system_resources, args=(5, 300))
    monitor_thread.start()
    
    # 运行压力测试
    tester = PressureTester()
    tester.run_pressure_test(total_requests=100, concurrent_workers=10)
    
    monitor_thread.join()

7. 总结

通过今天的实战,我们完成了从单次API调用到批量处理,再到压力测试的完整流程。让我们回顾一下关键点:

7.1 核心收获

  1. API调用基础:掌握了如何通过Python代码调用vLLM部署的Qwen2.5-72B-Instruct模型,而不仅仅是通过Chainlit界面。

  2. 批量处理能力:学会了两种批量处理方式:

    • 顺序调用:简单可靠,适合小批量任务
    • 并发调用:高效快速,适合大批量任务
  3. 压力测试实战:实现了完整的压力测试方案,能够:

    • 测试不同并发数下的性能表现
    • 分析响应时间分布和成功率
    • 发现性能瓶颈和优化方向
  4. 性能监控:学会了如何监控系统资源使用情况,为性能优化提供数据支持。

7.2 实际应用建议

根据我的经验,这里有一些实用建议:

  1. 并发数选择:对于72B参数的大模型,建议从较低的并发数(如3-5)开始测试,根据服务器配置逐步增加。一般来说,单GPU服务器建议并发数不要超过10。

  2. 请求频率控制:即使使用并发,也建议添加适当的请求间隔(如0.1-0.5秒),避免瞬间压力过大。

  3. 错误处理:一定要实现完善的错误处理和重试机制,特别是对于生产环境。

  4. 结果缓存:对于重复性较高的查询,可以考虑实现结果缓存,减少对模型的重复调用。

  5. 监控告警:建立监控系统,关注响应时间、错误率、Token使用量等关键指标。

7.3 下一步探索方向

如果你已经掌握了批量调用和压力测试,可以进一步探索:

  1. 异步调用:使用asyncioaiohttp实现真正的异步请求,进一步提高效率。

  2. 负载均衡:如果有多台服务器,可以实现负载均衡,分散请求压力。

  3. 请求优先级:为不同重要性的请求设置优先级,确保关键任务优先处理。

  4. 流式响应:对于长文本生成,使用流式响应改善用户体验。

  5. 成本优化:分析Token使用情况,优化提示词,降低使用成本。

Qwen2.5-72B-Instruct-GPTQ-Int4是一个功能强大的模型,通过合理的批量调用和并发控制,你可以充分发挥它的潜力。记住,性能优化是一个持续的过程,需要根据实际使用情况不断调整和优化。

希望这篇实战指南能帮助你在实际项目中更好地使用这个大模型。如果在使用过程中遇到问题,或者有更好的优化建议,欢迎交流分享。


获取更多AI镜像

想探索更多AI镜像和应用场景?访问 CSDN星图镜像广场,提供丰富的预置镜像,覆盖大模型推理、图像生成、视频生成、模型微调等多个领域,支持一键部署。

Logo

这里是“一人公司”的成长家园。我们提供从产品曝光、技术变现到法律财税的全栈内容,并连接云服务、办公空间等稀缺资源,助你专注创造,无忧运营。

更多推荐