基于asyncio优化process_coordinates函数实现4坐标并行请求
优化process_coordinates函数实现批量并行请求街景数据
假设你的原有函数大致结构如下(包含异步请求、数据转换和Parquet存储逻辑):
import asyncio import pandas as pd async def find_panorama_async(lat, lon): # 原有异步请求逻辑,返回单条坐标的街景数据 pass async def process_coordinates(coordinates): result_dict = {} # 原有逐个迭代请求的逻辑 for idx, (lat, lon) in enumerate(coordinates): data = await find_panorama_async(lat, lon) result_dict[idx] = data # 原有数据转换流程(示例) df = pd.DataFrame.from_dict(result_dict, orient='index') # 原有Parquet存储流程 df.to_parquet('panorama_data.parquet') return result_dict
优化后的实现(并发限制为4个请求)
我们可以用asyncio.Semaphore来控制并发量,确保同时最多处理4个坐标请求,同时保留原有数据转换和存储逻辑:
import asyncio import pandas as pd async def find_panorama_async(lat, lon): # 保留原有异步请求逻辑 pass async def _bound_request(sem, lat, lon): # 用信号量包装请求,控制并发数 async with sem: return await find_panorama_async(lat, lon) async def process_coordinates(coordinates): # 初始化信号量,限制并发数为4 sem = asyncio.Semaphore(4) result_dict = {} # 创建所有并发任务 tasks = [] for idx, (lat, lon) in enumerate(coordinates): # 绑定信号量和当前坐标任务,同时记录索引 task = asyncio.create_task(_bound_request(sem, lat, lon)) tasks.append((idx, task)) # 等待所有任务完成,收集结果到字典 for idx, task in tasks: result_dict[idx] = await task # 保留原有数据转换流程(此处沿用你的原有逻辑) df = pd.DataFrame.from_dict(result_dict, orient='index') # 保留原有Parquet存储流程 df.to_parquet('panorama_data.parquet') return result_dict
关键优化点说明
- 并发控制:使用
asyncio.Semaphore(4)确保同一时间只有4个异步请求在运行,避免因请求过载被服务端限制或触发反爬机制。 - 任务批量处理:将所有坐标请求打包成异步任务,通过
asyncio.create_task提交,最后用await逐个收集结果并映射到原索引,保证结果顺序与输入坐标一致。 - 原有逻辑保留:完全保留了结果字典的构建、数据转换(转DataFrame)和Parquet文件存储的流程,无需修改原有业务逻辑。
额外建议
如果需要处理请求失败的情况,可以在_bound_request中添加异常捕获,比如:
async def _bound_request(sem, lat, lon): async with sem: try: return await find_panorama_async(lat, lon) except Exception as e: print(f"请求坐标({lat}, {lon})失败: {str(e)}") return None # 或返回自定义错误标记
内容的提问来源于stack exchange,提问作者Daniel AG
相关产品推荐
相关产品推荐

