Python ThreadPool工作线程失败时更新全局变量的实现方案
解决方案与分析
一、逻辑实现:thread()函数 + 线程安全的共享状态
要实现“密钥过期时暂停所有线程、更新后继续”的逻辑,不能只在api_call里单独处理——因为需要跨线程同步暂停信号和密钥更新,必须结合线程安全的共享变量来控制。具体来说,需要在thread()函数所在的上下文里引入以下线程同步工具:
- 一个锁(
threading.Lock):确保只有一个线程去请求新密钥,避免重复操作 - 一个事件(
threading.Event):用来全局暂停/恢复所有工作线程 - 共享的请求头变量:让所有线程使用统一的最新密钥
修改后的完整代码如下:
import requests import json from multiprocessing.pool import ThreadPool from tqdm import tqdm import threading # 线程安全的共享状态初始化 auth_update_lock = threading.Lock() pause_workers = threading.Event() pause_workers.set() # 初始状态:允许所有线程执行 current_request_headers = {"authorization": get_new_authorization()} def api_call(result): global current_request_headers while True: # 等待暂停信号解除,没解除时线程会阻塞在这里 pause_workers.wait() try: resp = requests.get(result["url"], headers=current_request_headers) resp.raise_for_status() # 主动触发HTTP错误(比如401密钥过期) return json.loads(resp.text) except requests.exceptions.HTTPError as err: if err.response.status_code == 401: # 尝试更新密钥:加锁确保只有一个线程执行更新逻辑 with auth_update_lock: # 二次检查:避免多个线程同时触发更新时重复请求新密钥 if pause_workers.is_set(): pause_workers.clear() # 暂停所有线程 # 获取新密钥并更新共享头 current_request_headers["authorization"] = get_new_authorization() pause_workers.set() # 恢复所有线程继续执行 else: # 其他HTTP错误(如404、500),按需求处理,这里返回None忽略 return None except Exception as e: # 其他异常(如网络超时),同样按需求处理 return None def thread(worker, jobs, num_threads=5): pool = ThreadPool(num_threads) results = [] for result in tqdm(pool.imap_unordered(worker, jobs), total=len(jobs)): if result: results.append(result) pool.close() pool.join() return results if results else None def get_new_authorization(): # 这里替换成你的密钥获取逻辑 return "fresh_auth_token" # 使用示例 input_results = [...] # 你的输入数据列表 final_results = thread(api_call, input_results)
关键逻辑说明:
- 当某线程检测到401错误时,先通过锁抢占密钥更新权限,避免多线程重复请求新密钥
- 更新密钥前,通过
pause_workers.clear()让所有线程阻塞在pause_workers.wait()处 - 密钥更新完成后,用
pause_workers.set()恢复所有线程,此时所有线程会使用新的共享请求头继续请求
二、ThreadPool是否是最佳方案?
对于API请求这种IO密集型任务,ThreadPool是非常合适的选择:
- 线程的创建、切换开销远低于进程,资源占用更少
- Python的GIL(全局解释器锁)在IO操作(比如网络请求)时会自动释放,不会影响多线程的并行效率
如果需要更现代的API接口,可以替换为concurrent.futures.ThreadPoolExecutor(功能和ThreadPool一致,但支持上下文管理器,代码更简洁),但核心的线程同步逻辑完全通用。只有当你处理CPU密集型任务时,才需要考虑进程池——显然你的场景不适用。
内容的提问来源于stack exchange,提问作者OJT
相关产品推荐
相关产品推荐

