Pydantic Evals 并发控制与性能调优:掌握 max_concurrency 的完整实战指南
2026/9/14 2:26:04 网站建设 项目流程

Pydantic Evals 并发控制与性能调优:掌握 max_concurrency 的完整实战指南

【免费下载链接】pydantic-aiHow Python does AI. Agents, realtime voice, image generation, embeddings. Every model, every interface, typed end to end.项目地址: https://gitcode.com/GitHub_Trending/py/pydantic-ai

Pydantic Evals 是 pydantic-ai 仓库中用于离线评测(Dataset Evaluation)与在线评测(Online Evaluation)的框架,默认会以最大吞吐量的方式并发执行所有评测用例。本文以Dataset.evaluate()/Dataset.evaluate_sync()max_concurrency参数为核心,系统讲解并发执行的默认行为、限制并发的典型场景(API 限流、资源约束、调试)、同步/异步 API 的一致性、性能对比测算方法,并结合仓库源码(pydantic_evals/pydantic_evals/dataset.py、pydantic_evals/pydantic_evals/_utils.py)与测试用例(tests/evals/test_dataset.py)揭示底层实现原理。读完本文,你将能针对不同评测场景选择正确的并发度、估算有效并发率,并与重试策略配合构建稳定高效的评测流水线。

概览:默认全量并发

在 Pydantic Evals 中,一次评测(evaluation)由两部分组成:Task(被测函数,接收 case 的 inputs 并产生 output)与Evaluator(对 output 进行评判的打分器)。框架默认行为是将所有 case 并发执行以最大化吞吐量

这一默认行为在源码中有明确体现。在 dataset.py 中,evaluate方法的max_concurrency参数声明为int | None = None,其文档注释写明:

max_concurrency: The maximum number of concurrent evaluations of the task to allow. If None, all cases will be evaluated concurrently.

也就是说,不传该参数时,len(self.cases)个用例会同时进入执行;传了max_concurrency=N时,则同一时刻最多只有 N 个用例在执行,其余用例排队等待。

控制并发的方式是信号量(Semaphore)。源码第 337 行:

limiter = anyio.Semaphore(max_concurrency) if max_concurrency is not None else AsyncExitStack()
  • 传入max_concurrency时,创建anyio.Semaphore(N),每个 case 的执行协程通过async with limiter:(第 362 行)获取许可,从而限制同时运行的 case 数量;
  • 不传时退化为AsyncExitStack()(无限制上下文管理器),所有 case 立即并发运行。

所有 case 的协程通过 task_group_gather 在anyio.create_task_group()中并发调度,并按输入顺序收集结果。这意味着并发调度本身是异步 I/O 层面的协作式并发,适合 LLM API 调用、数据库查询、网络请求等场景。

基本用法:三个档位控制并发度

最基本的用法非常直观:evaluate_sync接受一个普通(同步)或异步 task 函数,通过max_concurrency指定并发档位。

from pydantic_evals import Case, Dataset def my_task(inputs: str) -> str: return f'Result: {inputs}' dataset = Dataset(name='concurrency_demo', cases=[Case(inputs='test1'), Case(inputs='test2')]) # 所有 case 并发执行(默认行为) report = dataset.evaluate_sync(my_task) # 限制为 5 个并发 case report = dataset.evaluate_sync(my_task, max_concurrency=5) # 串行执行(一次一个) report = dataset.evaluate_sync(my_task, max_concurrency=1)

三个档位的语义:

max_concurrency取值行为
不传 /None所有 case 全量并发,吞吐量最大
N(正整数)同时最多执行 N 个 case,其余排队
1严格串行,一次只处理一个 case

合法性校验:源码在 dataset.py 中显式校验:

if max_concurrency is not None and max_concurrency < 1: raise ValueError(f'max_concurrency must be >= 1, got {max_concurrency}')

传入0或负数会直接抛出ValueError。这一点同样被测试用例覆盖:tests/evals/test_dataset.py中的test_evaluate_with_invalid_max_concurrency[0, -1]参数化断言了该异常。

何时需要限制并发

默认全量并发虽然吞吐最大,但在真实场景中往往需要主动限流。以下是三个最典型的场景。

场景一:API 限流(Rate Limiting)

多数 LLM API 都会对并发请求数或每秒请求数(RPS)设限。如果你的服务允许每秒 10 个请求,就把max_concurrency设为 10:

from pydantic_evals import Case, Dataset async def my_llm_task(inputs: str) -> str: return f'LLM Result: {inputs}' dataset = Dataset(name='rate_limit_demo', cases=[Case(inputs='test1')]) # 如果你的 API 允许 10 个并发请求/秒 report = dataset.evaluate_sync( my_llm_task, max_concurrency=10, )

这里的值需要根据你实际使用的模型服务商(OpenAI、Anthropic、Google 等,见 docs/models)的限流文档来设定。若并发过高,评测会因限流异常而失败,此时除了降低并发,还应配合 重试策略 中的指数退避来吸收瞬时 429 错误。

场景二:资源约束(Resource Constraints)

当 task 或 evaluator 是内存密集、CPU 密集操作,或依赖有限的共享资源(如数据库连接池)时,过高的并发会耗尽系统资源:

from pydantic_evals import Case, Dataset def heavy_computation(inputs: str) -> str: return f'Heavy: {inputs}' def db_query_task(inputs: str) -> str: return f'DB: {inputs}' dataset = Dataset(name='resource_constraints', cases=[Case(inputs='test1')]) # 内存密集型操作:同一时刻只跑 2 个 report = dataset.evaluate_sync( heavy_computation, max_concurrency=2, # 一次最多 2 个 ) # 数据库连接池上限:与连接池大小对齐 report = dataset.evaluate_sync( db_query_task, max_concurrency=5, # 匹配连接池大小 )

最佳实践是把max_concurrency设为与底层资源容量相等或略低的值,例如数据库连接池大小为 5 就设 5,避免排队等待连接导致的超时堆积。

场景三:调试(Debugging)

并发执行时,多个 case 的异常交错出现,错误堆栈难以定位。调试阶段使用max_concurrency=1串行执行,可以获得清晰、可复现的错误轨迹:

from pydantic_evals import Case, Dataset def my_task(inputs: str) -> str: return f'Result: {inputs}' dataset = Dataset(name='debug_demo', cases=[Case(inputs='test1')]) # 更易于调试 report = dataset.evaluate_sync( my_task, max_concurrency=1, )

串行模式还保证了评测结果的可复现性,适合在 CI 中对失败用例做最小化复现。

性能对比:并发度如何影响总耗时

下面是一个可直接运行的对比示例,用 10 个模拟 100ms 延迟的 case 展示不同并发度下的耗时差异:

import asyncio from pydantic_evals import Case, Dataset # 构造包含多个测试用例的数据集 dataset = Dataset( name='performance_comparison', cases=[ Case( name=f'case_{i}', inputs=i, expected_output=i * 2, ) for i in range(10) ] ) async def slow_task(input_value: int) -> int: """模拟慢操作(如 API 调用)。""" await asyncio.sleep(0.1) # 每个 case 100ms return input_value * 2 # 无限制并发:总耗时约 0.1s(所有 case 并行执行) report = dataset.evaluate_sync(slow_task) # 限制并发 2:总耗时约 0.5s(每次 2 个,共 5 批) report = dataset.evaluate_sync(slow_task, max_concurrency=2) # 串行:总耗时约 1.0s(一次 1 个,共 10 个 case) report = dataset.evaluate_sync(slow_task, max_concurrency=1)

估算逻辑(假设每个 case 耗时 T、共 N 个 case、并发度为 C):理想总耗时 ≈ceil(N / C) × T。上例中 N=10、T=0.1s:

  • C=∞(默认):1 × 0.1s = 0.1s
  • C=2:5 × 0.1s = 0.5s
  • C=1:10 × 0.1s = 1.0s

需要注意的是,这是理想化估算。实际中并发带来的上下文切换、共享资源竞争、以及 task 内部的 CPU 密集计算,都会使加速比低于线性,特别是对CPU 密集而非 I/O 密集的任务,Pydantic Evals 的并发模型是 asyncio 协作式并发,不会真正利用多核。对于 CPU 密集任务,应结合thread_executor(Pydantic AI 的 并发能力)等方案另行处理。

与 Evaluator 的并发协作

默认情况下,Task 执行与 Evaluator 执行都是并发的。一个 case 的完整流程是:先运行 task 得到 output,再依次运行该 case 绑定的 evaluators(数据集级evaluators或 case 级evaluators),两者都在同一个_handle_case协程内、受同一个信号量限制。

from pydantic_evals import Case, Dataset from pydantic_evals.evaluators import LLMJudge def my_task(inputs: str) -> str: return f'Result: {inputs}' dataset = Dataset( name='evaluator_concurrency', cases=[Case(inputs=f'test{i}') for i in range(100)], # 100 个 case evaluators=[ LLMJudge(rubric='Quality check'), # 内部会发起 LLM API 调用 ], ) # task 与 evaluator 均在该并发度下受控运行 report = dataset.evaluate_sync( my_task, max_concurrency=10, )

如果 evaluator 本身很昂贵(例如LLMJudge这类基于 LLM 的裁判模型,定义于 pydantic_evals/pydantic_evals/evaluators/llm_as_a_judge.py),限制并发可以帮助管理:

  • API 限流:task 与 evaluator 都在调用模型 API,两者的并发请求会叠加,需要统一计入限流预算;
  • 成本:更少的并发 API 调用意味着更平稳的 token 消耗;
  • 内存占用:并发执行的 case 会同时持有中间结果与上下文,限制并发可控制峰值内存。

由于信号量覆盖整个 case 生命周期(从 task 到 evaluator 再到报告收集),max_concurrency=10意味着"同时处于运行中的 case 不超过 10 个",而不是"task 并发 10、evaluator 再并发 10"。

同步 API 与异步 API:行为完全一致

evaluate_syncevaluate在并发控制上行为完全一致,区别仅在于调用方式。

同步 API:evaluate_sync

from pydantic_evals import Case, Dataset def my_task(inputs: str) -> str: return f'Result: {inputs}' dataset = Dataset(name='sync_demo', cases=[Case(inputs='test1')]) # 内部通过事件循环运行异步任务,并应用受控并发 report = dataset.evaluate_sync(my_task, max_concurrency=10)

从源码看,evaluate_syncevaluate的薄包装:见 dataset.py,它把全部参数原样透传给evaluate,并用 run_until_complete 在事件循环上驱动协程直至完成——即使你传入的 task 是同步函数,也会在事件循环内部以兼容方式运行(_run_task同时接受同步与异步 callable)。因此同步 task 同样能享受 I/O 并发(如内部使用httpxrequests等阻塞调用时仍需谨慎,阻塞调用会卡住事件循环,建议在 async task 中使用异步客户端)。

异步 API:evaluate

from pydantic_evals import Case, Dataset async def my_task(inputs: str) -> str: return f'Result: {inputs}' async def run_evaluation(): dataset = Dataset(name='async_demo', cases=[Case(inputs='test1')]) # 行为与同步版一致,但运行在异步上下文中 report = await dataset.evaluate(my_task, max_concurrency=10) return report

两者共享同一套参数签名(namemax_concurrencyprogressretry_taskretry_evaluatorstask_namemetadatarepeatlifecycle),在大型评测管线中可以直接在异步服务内调用evaluate,避免阻塞事件循环。

监控并发:量化你的评测效率

仅凭直觉设置并发度是不够的,建议用耗时统计来量化"有效并发率",从而判断当前并发设置是否合理:

import time from pydantic_evals import Case, Dataset def task(inputs: str) -> str: return f'Result: {inputs}' dataset = Dataset(name='monitoring', cases=[Case(inputs=f'test{i}') for i in range(10)]) t0 = time.time() report = dataset.evaluate_sync(task, max_concurrency=10) duration = time.time() - t0 num_cases = len(report.cases) + len(report.failures) avg_duration = duration / num_cases print(f'Total: {duration:.2f}s') #> Total: 0.01s print(f'Cases: {num_cases}') #> Cases: 10 print(f'Avg per case: {avg_duration:.2f}s') #> Avg per case: 0.00s print(f'Effective concurrency: ~{num_cases * avg_duration / duration:.1f}') #> Effective concurrency: ~1.0

有效并发率公式有效并发率 ≈ (case 总数 × 单 case 平均耗时) / 总耗时。它反映的是"平均同时有多少个 case 在真正运行":

  • 若远低于max_concurrency,说明瓶颈不在并发度(可能是 task 内部串行、共享资源排队、或 CPU 密集无法并行);
  • 若接近max_concurrency,说明并发设置合理,增大并发可能进一步缩短总耗时;
  • 若总耗时增长但有效并发率不升,说明已触达资源或限流上限,继续加大并发只会增加错误率。

报告中还包含每个 case 的task_durationtotal_duration(见ReportCase字段),可以进一步区分 task 耗时与 evaluator 耗时,精确定位慢在哪一环。

处理限流:并发与重试的配合

如果评测过程中仍然命中限流,评测会失败。两个手段配合使用:

  1. 降低并发度,从源头减少瞬时请求量;
  2. 配置重试策略retry_task/retry_evaluators,基于 Tenacity),吸收少量瞬时失败。
from pydantic_evals import Case, Dataset def task(inputs: str) -> str: return f'Result: {inputs}' dataset = Dataset(name='rate_limit_handling', cases=[Case(inputs='test1')]) # 降低并发以规避限流 report = dataset.evaluate_sync( task, max_concurrency=5, # 保持在限流阈值之下 )

更完整的组合方案见 重试策略:例如stop_after_attempt(5)+wait_exponential(multiplier=2, min=2, max=60)的指数退避配置,配合max_concurrency双管齐下。该文档也给出了一条运维经验:"持续命中限流"时,应增大退避间隔(而非仅增加重试次数),并同步降低max_concurrency

更多并发相关参数

evaluate/evaluate_sync的完整签名(见 dataset.py)中,以下参数与并发性能密切相关:

参数类型默认值说明
max_concurrencyint \| NoneNone最大并发 case 数;None表示全量并发,< 1ValueError
progressboolTrue是否显示评测进度条,进度条按完成 case 数推进
retry_taskRetryConfig \| NoneNonetask 执行失败时的重试配置(Tenacity)
retry_evaluatorsRetryConfig \| NoneNoneevaluator 执行失败时的重试配置(Tenacity)
repeatint1每个 case 重复执行的次数;> 1时结果按原 case 名分组聚合
lifecycleCaseLifecycleNone每个 case 的 setup / prepare_context / teardown 钩子,在信号量内按 case 实例化

关于repeat的一个并发细节:源码_build_tasks_to_run(dataset.py)会把每个 case 展开为repeat个独立任务(命名如case [1/3]),这些任务同样受max_concurrency限制。因此若需评估同一输入的多次运行以统计方差,请合理设定并发度,避免重复运行把限流风险放大repeat倍。更细的用法可参考 多轮运行指南。

大数据集场景

面对大规模数据集(数千个 case),建议:

  • 先小规模标定:用少量 case + 时间统计确定单 case 耗时与限流预算,再推算合适的max_concurrency
  • 使用数据集管理:数据集管理指南 介绍了Dataset.from_file/to_file等序列化能力(YAML/JSON),可把大数据集拆分到文件、按需加载;
  • 配合观测:Logfire 集成 可以在追踪中查看每个 case 的耗时、重试次数与失败分布,帮助持续调优并发参数。

源码级验证:测试如何保障并发行为

仓库测试对并发控制有直接覆盖,可作为行为契约的参考(tests/evals/test_dataset.py):

  • test_evaluate_with_concurrency:用max_concurrency=1串行执行包含 2 个 case 的数据集,并完整断言报告的字段(assertions、scores、duration、span/trace id 等),验证受限并发下结果正确性与报告完整性;
  • test_evaluate_with_invalid_max_concurrency:参数化[0, -1],断言抛出ValueError('max_concurrency must be >= 1, got ...')

这两组测试与 dataset.py 的参数校验、anyio.Semaphore限制逻辑一一对应,证明并发控制是框架的正式、受测试保障的能力,而非文档中的口头承诺。

下一步

  • 重试策略:处理瞬时失败,与并发控制配合应对限流
  • 数据集管理:处理大型数据集,序列化与按需加载
  • Logfire 集成:观测评测性能、重试与失败分布

【免费下载链接】pydantic-aiHow Python does AI. Agents, realtime voice, image generation, embeddings. Every model, every interface, typed end to end.项目地址: https://gitcode.com/GitHub_Trending/py/pydantic-ai

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询