并行处理HTTPS请求丢失问题求助(基于Questrade API批量调用)
解决Questrade API并行请求丢失的问题
首先,我注意到你当前的实现可能是每个股票代码单独发起请求,但结合Questrade API的批量查询能力(支持逗号分隔的ID列表),先把11000+代码分成每100个一组再发起请求,能大幅减少总请求数,从根源降低请求丢失的概率。再配合以下优化手段,就能彻底解决HTTPS请求丢失的问题:
1. 正确实现批量请求(核心优化)
Questrade的markets/quotes/{ids}接口支持一次性传入最多100个逗号分隔的股票ID,所以先把你的代码列表按100个一组拆分,每组只发一次请求,而不是每个ID单独请求。这样总请求数从11000+降到110左右,并发压力骤减。
示例代码:
import requests from joblib import Parallel, delayed from itertools import islice def chunk_iterable(iterable, chunk_size): """把可迭代对象按指定大小拆分批次""" it = iter(iterable) while True: chunk = tuple(islice(it, chunk_size)) if not chunk: break yield chunk def batch_request(self, id_chunk, base_url, target_key): # 把批次内的ID拼成逗号分隔的字符串 ids_str = ','.join(id_chunk) full_url = f"{base_url}{ids_str}" try: response = requests.get(full_url, headers=self.headers, timeout=10) response.raise_for_status() # 主动触发HTTP错误异常 return response.json().get(target_key, []) except requests.exceptions.RequestException as e: print(f"批次请求失败: {full_url}, 错误详情: {str(e)}") return [] # 使用方式 stock_ids = [你的11000+股票代码列表] chunked_ids = chunk_iterable(stock_ids, 100) # 每100个ID为一组 batch_results = Parallel(n_jobs=10, verbose=10)( delayed(self.batch_request)(chunk, api_base_url, "quotes") for chunk in chunked_ids ) # 合并所有批次的结果 final_data = [] for res in batch_results: final_data.extend(res)
2. 控制并发数,避免API过载
Parallel(n_jobs=-1)会占用所有CPU核心,导致并发请求数过高,触发Questrade的API限流或本地连接池耗尽。建议设置合理的并发数,比如n_jobs=10或n_jobs=20,可以根据Questrade API的速率限制文档调整这个数值。
3. 添加重试机制,处理临时网络问题
网络波动、API临时限流都可能导致请求失败,添加自动重试能有效解决这类偶发问题。可以用tenacity库实现指数退避重试:
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type @retry( stop=stop_after_attempt(3), # 最多重试3次 wait=wait_exponential(multiplier=1, min=2, max=10), # 指数退避等待(2s→4s→8s) retry=retry_if_exception_type((requests.exceptions.ConnectionError, requests.exceptions.HTTPError)) ) def batch_request(self, id_chunk, base_url, target_key): ids_str = ','.join(id_chunk) full_url = f"{base_url}{ids_str}" response = requests.get(full_url, headers=self.headers, timeout=10) response.raise_for_status() return response.json().get(target_key, [])
4. 使用Session复用连接
requests.Session会复用TCP连接,减少握手开销,避免因频繁创建连接导致的资源耗尽。可以在类的初始化方法中创建Session:
def __init__(self): self.headers = {"Authorization": "Bearer 你的API密钥"} self.session = requests.Session() self.session.headers.update(self.headers) @retry(...) def batch_request(self, id_chunk, base_url, target_key): ids_str = ','.join(id_chunk) full_url = f"{base_url}{ids_str}" response = self.session.get(full_url, timeout=10) response.raise_for_status() return response.json().get(target_key, [])
5. 增加错误日志与单独重试机制
在捕获请求异常时,记录失败的批次ID,后续可以单独对这些批次发起重试,避免遗漏数据。比如把失败的批次存入列表,最后统一处理:
failed_chunks = [] def batch_request(self, id_chunk, base_url, target_key): ids_str = ','.join(id_chunk) full_url = f"{base_url}{ids_str}" try: response = self.session.get(full_url, timeout=10) response.raise_for_status() return response.json().get(target_key, []) except requests.exceptions.RequestException as e: print(f"批次请求失败: {full_url}, 错误详情: {str(e)}") failed_chunks.append(id_chunk) return [] # 批量请求完成后,重试失败的批次 if failed_chunks: retry_results = Parallel(n_jobs=5)( delayed(self.batch_request)(chunk, api_base_url, "quotes") for chunk in failed_chunks ) for res in retry_results: final_data.extend(res)
内容的提问来源于stack exchange,提问作者Jeremie
相关产品推荐
相关产品推荐

