You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何控制Python中ThreadPoolExecutor的任务吞吐量速度?

如何控制ThreadPoolExecutor的吞吐量速度?

嘿,你现在用concurrent.futures.ThreadPoolExecutor跑异步任务还配了tqdm监控进度,想控制吞吐量是吧?这需求太常见了——尤其是当你爬的网站有访问频率限制,或者不想一下子把服务器资源占满的时候。我给你几个实用的方法,你按需选:

方法1:用信号量(Semaphore)精准控制并发执行数

虽然ThreadPoolExecutor的max_workers能控制线程总数,但如果想更细粒度地限制同时执行的任务数量(比如线程池开了10个线程,但只想同时跑5个任务),threading.Semaphore是绝佳选择。它就像一个“许可证池”,每个任务执行前必须先拿到许可证,执行完再归还,从而保证同一时间最多有指定数量的任务在运行。

代码示例:

import threading
import concurrent.futures
from tqdm import tqdm

# 定义信号量,设置同时允许执行的任务数,比如设为5
semaphore = threading.Semaphore(5)

def target_function(URL):
    # 任务执行前先获取信号量许可
    with semaphore:
        # 这里写你原有的任务逻辑,比如请求URL、处理响应等
        # ... 你的代码 ...
        pass

n_jobs = 10
URL_list = ["url1", "url2", "url3", ...]  # 你的URL列表

with concurrent.futures.ThreadPoolExecutor(max_workers=n_jobs) as executor:
    future_to_url = {executor.submit(target_function, URL): URL for URL in URL_list}
    # 用tqdm监控任务完成进度
    for future in tqdm(concurrent.futures.as_completed(future_to_url),
                      total=len(future_to_url),
                      unit='URL',
                      unit_scale=True,
                      leave=False):
        url = future_to_url[future]
        try:
            result = future.result()
            # 处理任务结果
        except Exception as exc:
            print(f"处理{url}时发生异常: {exc}")

这种方法侵入性低,只需要给目标函数加个with semaphore块,就能灵活控制吞吐量,完全不影响你原有的tqdm监控逻辑。

方法2:分批提交任务,控制提交节奏

如果你想严格控制任务的提交速率(比如每秒提交5个,或者每完成一批再提交下一批),可以把任务分成批次,提交一批后等待部分任务完成,再继续提交下一批。

代码示例:

import concurrent.futures
from tqdm import tqdm
import time

n_jobs = 10
batch_size = 5  # 每批提交5个任务
URL_list = ["url1", "url2", "url3", ...]
total_urls = len(URL_list)

with concurrent.futures.ThreadPoolExecutor(max_workers=n_jobs) as executor:
    futures = []
    pbar = tqdm(total=total_urls, unit='URL', unit_scale=True, leave=False)
    
    # 按批次提交任务
    for i in range(0, total_urls, batch_size):
        current_batch = URL_list[i:i+batch_size]
        # 提交当前批次的任务
        for url in current_batch:
            future = executor.submit(target_function, url)
            futures.append((future, url))
        
        # 等待当前批次至少完成1个任务(可根据需求调整等待逻辑)
        completed_in_batch = 0
        for future, url in futures[-len(current_batch):]:
            if future.done():
                completed_in_batch += 1
                try:
                    result = future.result()
                    # 处理结果
                except Exception as exc:
                    print(f"处理{url}时发生异常: {exc}")
                pbar.update(1)
        
        # 可选:加个固定延迟,控制批次提交间隔
        time.sleep(1)
    
    # 处理剩余未完成的任务
    for future, url in futures:
        if not future.done():
            try:
                result = future.result()
                # 处理结果
            except Exception as exc:
                print(f"处理{url}时发生异常: {exc}")
            pbar.update(1)
    pbar.close()

这种方法适合你需要严格把控任务提交节奏的场景,避免瞬间把所有任务都扔进线程池。

方法3:给任务加固定延迟(简单粗暴版)

如果场景比较简单,你可以直接在目标函数里加固定延迟,让每个任务执行完后暂停一会儿,自然降低整体吞吐量。比如:

import time

def target_function(URL):
    # 执行你的任务逻辑
    # ... 你的代码 ...
    
    # 延迟1秒,控制每秒最多处理N个任务(N等于线程池的max_workers)
    time.sleep(1)

这种方法虽然简单,但灵活性差——如果线程池开了10个线程,那每秒最多能处理10个任务,没法单独控制吞吐量和线程数的关系,适合对精度要求不高的场景。


内容的提问来源于stack exchange,提问作者sudonym

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.22 07:42:42