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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 02:44:53