如何快速调用200万次OSRM API并将结果写入pandas DataFrame
异步批量处理OSRM请求实现方案
核心注意事项
- 公共OSRM服务有严格的请求频率限制,并发数建议初始设置为50-100,过高会被临时封禁IP,拖慢整体进度;如果是本地私有部署的OSRM服务,并发可调整到500-1000
- 提前过滤坐标为空的无效行,无需发起请求,直接赋值
None即可 - 200万行数据不要一次性生成所有异步任务,按1-2万条每批分片处理,避免内存占用过高
- 所有请求必须绑定原始行索引,异步返回顺序不固定,靠索引匹配结果才不会错位
- 增加超时、连接异常捕获,单条请求失败不中断整体任务
完整实现代码
import pandas as pd import asyncio import aiohttp import json from tqdm import tqdm # --------------- 配置项 --------------- CONCURRENCY = 80 # 并发数,公共服务别调太高 BATCH_SIZE = 10000 # 每批处理行数 TIMEOUT = 10 # 单条请求超时时间(秒) OSRM_URL_TPL = "http://router.project-osrm.org/route/v1/driving/{from_lng},{from_lat};{to_lng},{to_lat}?overview=false" # 注意OSRM接口坐标顺序是 经度,纬度,别写反! # 异步请求逻辑 async def fetch(row_idx, url, session, sem): async with sem: try: async with session.get(url, timeout=TIMEOUT) as resp: if resp.status != 200: return (row_idx, None, None) data = await resp.json() # 接口返回code正常才取结果 if data.get("code") == "Ok" and len(data.get("routes", [])) > 0: dist = data["routes"][0]["distance"] dur = data["routes"][0]["duration"] return (row_idx, dist, dur) else: return (row_idx, None, None) except Exception as e: # 所有异常(超时、连接断开、解析错误等)都返回空值 return (row_idx, None, None) async def process_batch(df_batch): sem = asyncio.Semaphore(CONCURRENCY) tasks = [] # 先初始化两列空值 df_batch[["distance", "duration"]] = None async with aiohttp.ClientSession() as session: for idx, row in df_batch.iterrows(): # 跳过坐标为空的行 if pd.isna(row["fromLat"]) or pd.isna(row["fromLong"]) or pd.isna(row["toLat"]) or pd.isna(row["toLong"]): continue # 注意OSRM要求坐标顺序是 经度,纬度,和表中纬度在前的存储顺序相反 url = OSRM_URL_TPL.format( from_lng=row["fromLong"], from_lat=row["fromLat"], to_lng=row["toLong"], to_lat=row["toLat"] ) tasks.append(fetch(idx, url, session, sem)) # 加进度条实时查看处理速度 for future in tqdm(asyncio.as_completed(tasks), total=len(tasks), desc="处理批次"): row_idx, dist, dur = await future df_batch.loc[row_idx, "distance"] = dist df_batch.loc[row_idx, "duration"] = dur return df_batch def main(df): # Windows系统运行需要放开下面这行注释 # asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy()) # 分批处理 result_frames = [] total_batch = (len(df) // BATCH_SIZE) + 1 for i in range(total_batch): print(f"开始处理第{i+1}/{total_batch}批") batch = df.iloc[i*BATCH_SIZE : (i+1)*BATCH_SIZE].copy() batch_res = asyncio.run(process_batch(batch)) result_frames.append(batch_res) return pd.concat(result_frames, ignore_index=False) # 测试用例 if __name__ == "__main__": d = { "customerId": [1, 2, 2, 3, 3], "fromLat": ["48.103190","48.541160", "25.907100", "41.948767", None], "fromLong": ["11.601773", "2.628249", "-100.338113", "12.742488", None], "toLat": ["53.446762", "41.374426", "23.994722", "49.316171", "41.948767"], "toLong": ["10.061031", "2.037211", "-104.754147", "-33.968543", "12.742488"], } df = pd.DataFrame(data=d) df = main(df) print(df)
关键坑点说明
- OSRM接口的坐标顺序是
经度,纬度,和表中存储的纬度在前、经度在后顺序相反,写反会返回路径不存在或者结果完全错误 - 不要直接用
asyncio.gather一次性等待所有任务,用asyncio.as_completed配合进度条可以实时看处理速度,遇到报错也能及时定位 - 公共OSRM服务不适合跑200万级别的批量请求,大概率会被限流,建议本地用docker部署OSRM服务,部署完成后只需要把代码里的URL改成
http://127.0.0.1:5000/route/v1/driving/...即可,处理速度能提升几十上百倍 - 如果需要更高稳定性,可以给fetch函数加2-3次重试逻辑,对503、超时的请求重新发起,减少空值占比
内容的提问来源于stack exchange,提问作者Test
相关产品推荐
相关产品推荐

