并行处理20K API请求并将结果合并为单个DataFrame
解决方案
核心实现逻辑
并行调用时,让每个线程完成「API调用→CSV解析→生成小DataFrame」的完整流程,最后把所有线程返回的小DataFrame合并即可。同时要严格控制并发数和调用速率,避免触发API限流。
代码示例
import pandas as pd from concurrent.futures import ThreadPoolExecutor, as_completed import time # 定义API调用与数据处理函数 def fetch_and_process_api(api_param): # 替换为你的实际API调用逻辑,比如用requests发起请求 # response = requests.get(f"your_api_url?param={api_param}") # 解析返回的CSV为DataFrame(示例用模拟数据,实际替换为response.content) csv_content = "col1,col2\nval_a,val_b\nval_c,val_d" df_chunk = pd.read_csv(pd.io.common.StringIO(csv_content)) # 在这里添加你的数据清洗、转换等处理逻辑 return df_chunk # 主执行逻辑 def main(): # 模拟20000个API请求参数(替换为你的实际参数列表) api_params = list(range(20000)) max_concurrent = 100 # 最大并发数 rate_limit_per_min = 1000 call_interval = 60 / rate_limit_per_min # 每次调用的最小间隔 df_chunks = [] call_counter = 0 start_timestamp = time.time() with ThreadPoolExecutor(max_workers=max_concurrent) as executor: # 提交所有任务到线程池 task_futures = {executor.submit(fetch_and_process_api, param): param for param in api_params} for future in as_completed(task_futures): try: chunk = future.result() df_chunks.append(chunk) except Exception as e: # 处理API调用失败的情况,比如记录日志、跳过失败任务 print(f"参数 {task_futures[future]} 调用失败: {str(e)}") # 速率控制:确保每分钟调用不超过1000次 call_counter += 1 elapsed_time = time.time() - start_timestamp expected_time = call_counter * call_interval if expected_time > elapsed_time: time.sleep(expected_time - elapsed_time) # 合并所有小DataFrame为最终结果 final_df = pd.concat(df_chunks, ignore_index=True) print(f"合并完成,总数据行数: {len(final_df)}") # 可选:保存到文件 # final_df.to_csv("merged_result.csv", index=False) if __name__ == "__main__": main()
关键注意事项
- 异常处理:必须捕获API调用可能出现的超时、返回非CSV数据等错误,避免单个任务失败导致整个流程中断。
- 内存优化:如果20000个小DataFrame占用内存过高,可以分批合并,比如每收集1000个就合并一次并清空临时列表。
- 速率控制:示例用简单的计数器和sleep实现限流,也可以用令牌桶算法优化,但这个场景下简单实现足够满足需求。
- 参数传递:确保
api_params是你的实际请求参数列表(比如不同的ID、查询条件等)。
内容的提问来源于stack exchange,提问作者Guillem Servera
相关产品推荐
相关产品推荐

