API不支持批量操作,如何加速150万条数据的邮箱查询?
优化150万条记录API调用效率的方案
当前你的多线程方案已经利用了部分并发能力,但还有几个关键优化点可以进一步缩短处理时间:
1. 避免全量加载CSV到内存
现有代码把150万行全部读入rows列表,会占用大量内存甚至触发OOM。改为逐行提交任务+批量写入,降低内存压力同时减少IO开销:
def process_records(sailthru_client, input_file, output_file, num_threads=600): start_time = time.time() with open(input_file, "r", encoding='utf-8') as csv_file, \ open(output_file, "w", newline="", encoding='utf-8') as output_csv: csv_reader = csv.DictReader(csv_file) fieldnames = csv_reader.fieldnames + ["email_address"] writer = csv.DictWriter(output_csv, fieldnames=fieldnames, quotechar='"', quoting=csv.QUOTE_ALL) writer.writeheader() with concurrent.futures.ThreadPoolExecutor(max_workers=num_threads) as executor: futures = [] # 逐行提交任务,避免内存过载 for row in csv_reader: future = executor.submit(convert_user_email_for_multithread, sailthru_client, row["Profile Id"]) futures.append((future, row)) # 批量处理结果,每1000条写入一次 batch = [] for future, row in futures: user_id, email_address = future.result() if email_address: row["email_address"] = email_address batch.append(row) else: print(f"Couldn't find email for user_id {user_id}. Skipping this row.") if len(batch) >= 1000: writer.writerows(batch) batch = [] # 写入剩余行 if batch: writer.writerows(batch) total_time = time.time() - start_time print(f"Total Processing Time: {total_time} seconds")
2. 优化速率控制逻辑
现有代码同时用Semaphore(300)和time.sleep(1/300),导致速率被双重限制,实际请求量低于300次/秒。改用令牌桶算法实现精准速率控制:
import time from threading import Lock class TokenBucket: def __init__(self, rate, capacity): self.rate = rate self.capacity = capacity self.tokens = capacity self.last_refill = time.time() self.lock = Lock() def acquire(self): with self.lock: now = time.time() # 按时间补充令牌 elapsed = now - self.last_refill new_tokens = elapsed * self.rate self.tokens = min(self.capacity, self.tokens + new_tokens) self.last_refill = now if self.tokens >= 1: self.tokens -= 1 return True else: # 等待直到有可用令牌 time.sleep((1 - self.tokens)/self.rate) self.tokens = 0 return True # 初始化令牌桶:300次/秒,容量300 token_bucket = TokenBucket(300, 300) def convert_user_email_for_multithread(sailthru_client, user_id): token_bucket.acquire() # 替代原有的semaphore和sleep try: response = sailthru_client.api_get("user", {"id": user_id}) if response.is_ok(): body = response.get_body() return user_id, body['keys']['email'] else: error = response.get_error() print(f"Error for user_id {user_id}: {error.get_message()}") print("Status Code:", response.get_status_code()) print("Error Code:", error.get_error_code()) return user_id, None except SailthruClientError as e: print(f"Exception for user_id {user_id}: {e}") return user_id, None
3. 调整线程池大小
API调用的瓶颈是网络IO而非CPU,线程池大小应设为速率限制的2-3倍(比如600-900),让更多线程等待网络响应,充分利用API的速率配额。现有代码num_threads=100远低于最优值,会导致请求无法达到300次/秒上限。
4. 增加重试机制
针对API临时错误(5xx状态码、超时等)添加重试逻辑,避免浪费请求配额:
import tenacity @tenacity.retry( stop=tenacity.stop_after_attempt(3), wait=tenacity.wait_exponential(multiplier=1, min=2, max=10), retry=tenacity.retry_if_exception_type((SailthruClientError,)) | tenacity.retry_if_result(lambda resp: resp and not resp.is_ok()) ) def fetch_user_email(sailthru_client, user_id): response = sailthru_client.api_get("user", {"id": user_id}) if response.is_ok(): return response.get_body()['keys']['email'] else: raise SailthruClientError(f"API error: {response.get_error().get_message()}") def convert_user_email_for_multithread(sailthru_client, user_id): token_bucket.acquire() try: email = fetch_user_email(sailthru_client, user_id) return user_id, email except SailthruClientError as e: print(f"Failed after retries for user_id {user_id}: {e}") return user_id, None
5. 切换到异步IO(推荐)
Python多线程受GIL限制,异步IO(asyncio+aiohttp)能更高效处理大量网络请求。如果Sailthru无官方异步客户端,可直接用aiohttp调用API:
import asyncio import aiohttp import csv import time async def fetch_user_email(session, api_key, api_secret, user_id): url = "https://api.sailthru.com/user" params = {"id": user_id, "api_key": api_key} auth = aiohttp.BasicAuth(api_key, api_secret) async with session.get(url, params=params, auth=auth) as response: if response.status == 200: data = await response.json() return user_id, data['keys']['email'] else: print(f"Error for user_id {user_id}: {await response.text()}") return user_id, None async def process_records_async(api_key, api_secret, input_file, output_file): start_time = time.time() semaphore = asyncio.Semaphore(300) async def bounded_fetch(user_id): async with semaphore: result = await fetch_user_email(session, api_key, api_secret, user_id) await asyncio.sleep(1/300) return result async with aiohttp.ClientSession() as session: with open(input_file, "r", encoding='utf-8') as csv_file, \ open(output_file, "w", newline="", encoding='utf-8') as output_csv: csv_reader = csv.DictReader(csv_file) fieldnames = csv_reader.fieldnames + ["email_address"] writer = csv.DictWriter(output_csv, fieldnames=fieldnames, quotechar='"', quoting=csv.QUOTE_ALL) writer.writeheader() tasks = [] rows = [] for row in csv_reader: rows.append(row) tasks.append(bounded_fetch(row["Profile Id"])) batch = [] for idx, result in enumerate(asyncio.as_completed(tasks)): user_id, email_address = await result row = rows[idx] if email_address: row["email_address"] = email_address batch.append(row) else: print(f"Couldn't find email for user_id {user_id}. Skipping this row.") if len(batch) >= 1000: writer.writerows(batch) batch = [] if batch: writer.writerows(batch) total_time = time.time() - start_time print(f"Total Processing Time: {total_time} seconds") # 执行异步任务 asyncio.run(process_records_async(api_key, api_secret, input_file, output_file))
预期效果
通过以上优化,尤其是异步IO+精准速率控制,处理时间可压缩至接近理论最小值(150万/300 = 5000秒 ≈ 1小时23分),实际因网络延迟会略有增加,但远低于当前的3-4小时。
内容的提问来源于stack exchange,提问作者Olivia Xu
相关产品推荐
相关产品推荐

