如何并行运行两个存在依赖关系的API?
并行处理依赖API请求的两个循环实现方案
需求是并行运行两个逻辑:第一个循环通过经纬度调用Geocode API生成place_id,第二个循环无需等待第一个循环全部完成,只要有place_id生成就立即调用Place Details API获取详情。原串行执行的方式效率极低,以下是基于Python线程池的优化实现:
核心思路
- 用线程安全队列传递
place_id:生产者线程(生成place_id)将结果存入队列,消费者线程(获取详情)从队列取数据处理,实现流式并行 - 用
concurrent.futures.ThreadPoolExecutor管理线程:IO密集型的API请求适合用线程池,避免多进程的开销 - 避免全局变量:改用线程安全的容器存储最终结果,防止并发写入冲突
优化后代码
import time import pandas as pd import requests import json from concurrent.futures import ThreadPoolExecutor, wait, ALL_COMPLETED from queue import Queue from typing import List, Dict API_KEY = "你的API密钥" # 控制并发数,根据Google API配额调整 PRODUCER_WORKERS = 5 CONSUMER_WORKERS = 5 # 线程安全队列:传递place_id place_id_queue = Queue(maxsize=20) # 线程安全结果容器 results: List[Dict] = [] def load_data() -> pd.DataFrame: """加载经纬度数据""" return pd.read_csv('lat_long file') def get_place_id(lat: float, lon: float) -> str | None: """调用Geocode API获取place_id""" url = f"https://maps.googleapis.com/maps/api/geocode/json?latlng={lat},{lon}&key={API_KEY}" try: response = requests.get(url, timeout=10) response.raise_for_status() data = response.json() if data['results']: place_id = data['results'][0]['place_id'] return place_id else: print(f"无结果:{lat},{lon}") return None except Exception as e: print(f"获取place_id失败 {lat},{lon}: {str(e)}") return None def producer_task(lat: float, lon: float): """生产者任务:生成place_id并存入队列""" place_id = get_place_id(lat, lon) if place_id: place_id_queue.put(place_id) def get_place_details(place_id: str) -> Dict | None: """调用Place Details API获取详情""" url = f"https://maps.googleapis.com/maps/api/place/details/json?language=en-US&placeid={place_id}&key={API_KEY}" try: response = requests.get(url, timeout=10) response.raise_for_status() data = response.json() result = data.get('result', {}) # 整理需要的字段 detail = { 'place_id': place_id, 'business_status': result.get('business_status', '0'), 'unit': '0' } # 提取premise信息 if result.get('address_components'): for component in result['address_components']: if component['types'][0] == 'premise': detail['unit'] = component['long_name'][4:] break return detail except Exception as e: print(f"获取详情失败 {place_id}: {str(e)}") return None def consumer_task(): """消费者任务:从队列取place_id并处理""" while True: place_id = place_id_queue.get() if place_id is None: # 终止信号 place_id_queue.task_done() break detail = get_place_details(place_id) if detail: results.append(detail) place_id_queue.task_done() def main(): df = load_data() lat_list = df['latitude'].tolist() lon_list = df['longitude'].tolist() # 启动消费者线程 consumer_executor = ThreadPoolExecutor(max_workers=CONSUMER_WORKERS) consumer_futures = [consumer_executor.submit(consumer_task) for _ in range(CONSUMER_WORKERS)] # 启动生产者线程,处理所有经纬度 producer_executor = ThreadPoolExecutor(max_workers=PRODUCER_WORKERS) producer_futures = [producer_executor.submit(producer_task, lat, lon) for lat, lon in zip(lat_list, lon_list)] # 等待所有生产者任务完成 wait(producer_futures, return_when=ALL_COMPLETED) # 向队列发送终止信号,每个消费者一个 for _ in range(CONSUMER_WORKERS): place_id_queue.put(None) # 等待所有消费者任务完成 wait(consumer_futures, return_when=ALL_COMPLETED) # 将结果转为DataFrame方便后续处理 result_df = pd.DataFrame(results) print(result_df.head()) # 保存结果 result_df.to_csv('place_details.csv', index=False) if __name__ == "__main__": main()
关键说明
- 队列控流:队列设置了
maxsize=20,防止生产者过快生成数据导致内存占用过高 - 并发数控制:
PRODUCER_WORKERS和CONSUMER_WORKERS可根据Google API的配额限制调整,避免触发限流 - 异常处理:每个API请求都添加了超时和状态码检查,错误信息更清晰
- 结果整理:用字典存储单条结果,最后转为DataFrame,比原代码的多个全局列表更易维护
- 优雅终止:生产者完成后向队列发送
None作为终止信号,确保消费者线程能正常退出
内容的提问来源于stack exchange,提问作者Hisham Ibrahim
相关产品推荐
相关产品推荐

