Pandas DataFrame图边列表生成优化:700万行性能瓶颈求解
我有一个以Geohash为索引的Pandas DataFrame,包含存储每个Geohash邻居列表的neighbours列,以及浪高、归一化浪高、Speed Factor等元数据列。需要生成用于图创建的edge_list,最初采用循环遍历的方式实现:
def create_edge_list(geohash, speed_factor, neighbours): edge_list = [] for n in neighbours: distance = haversine_distance(geohash, n) # distance is in km, speed is in m/s. speed = 14 * speed_factor time = round((distance/(speed*3.6))*60, 1) edge_list.append((geohash, n, {"distance": distance, "time": time})) return edge_list for geohash, row in tqdm(df.iterrows(), desc="Creating edge list", total=len(df.index), colour="green"): edge_list = create_edge_list(geohash, row.speed_factor, row.neighbours) elist.extend(edge_list)
但面对700万行数据时,这个方案速度极慢。随后尝试用ProcessPoolExecutor和ThreadPoolExecutor做多进程、多线程优化,起初效果不好,修复ProcessPoolExecutor的错误后,运行时间从数小时缩短到80分钟,但仍希望获得更多优化建议。以下是可复现代码:
# 使用Python 3.11.2,其他较新版本Python也可运行 !pip install geopandas !pip install geohash !pip install polygeohasher !pip install shapely !pip install pandas !pip install geopandas !pip install tqdm import os import random from math import cos, sin, asin, sqrt, radians import geohash as gh from polygeohasher import polygeohasher from shapely.wkt import loads import pandas as pd import geopandas as gpd from tqdm import tqdm def haversine_distance(geohash1, geohash2): # geohash2可能是邻居列表 if isinstance(geohash2, list): return [round(haversine_distance(geohash1, gh_val), 3) for gh_val in geohash2] lat1, lon1 = gh.decode(geohash1) lat2, lon2 = gh.decode(geohash2) lat1, lon1 = (float(lat1), float(lon1)) lat2, lon2 = (float(lat2), float(lon2)) lon1, lat1, lon2, lat2 = map(radians, [lon1, lat1, lon2, lat2]) # 哈弗辛公式 dlon = lon2 - lon1 dlat = lat2 - lat1 a = sin(dlat/2)**2 + cos(lat1) * cos(lat2) * sin(dlon/2)**2 c = 2 * asin(sqrt(a)) r = 6371 # 地球半径,单位公里。用3956则为英里,决定返回值单位。 return c * r def create_edge_list(geohash, speed_factor, neighbours): speed_multiplier = 60 / (3.6 * 14 * speed_factor) neighbours = list(neighbours) distances = haversine_distance(geohash, neighbours) times = [round(d * speed_multiplier, 2) for d in distances] edge_list = [(geohash, neighbours[i], {"distance": distances[i], "time": times[i]}) for i in range(len(times))] return edge_list if __name__ == "__main__": GEOHASH_PRECISION = 6 # 可通过工具创建多边形 polygon_wkt = "POLYGON((9.07196044921875 53.91728101547625,8.25897216796875 52.99495027026802,5.88043212890625 53.20603255157843,5.072937011718749 53.497849543967675,5.913391113281249 53.74221377343122,6.05621337890625 54.004540438503625,8.73687744140625 54.072282655603885,9.07196044921875 53.91728101547625))" polygon_gdf = gpd.GeoDataFrame(index=[0], crs="EPSG:4326", geometry=[loads(polygon_wkt)]) print("生成Geohash列表...") temp_df = polygeohasher.create_geohash_list(polygon_gdf, GEOHASH_PRECISION, inner=True) df = pd.DataFrame(temp_df.geohash_list.values.tolist()[0], columns=["geohash"]) df.set_index("geohash", inplace=True) # 模拟speed_factor数据 df["speed_factor"] = [random.uniform(0.4, 1.0) for i in range(len(df.index))] neighbours = {geohash: gh.neighbors(geohash) for geohash in df.index} df["neighbours"] = df.index.map(neighbours) elist = [] MT = False print("生成Edge List...") if MT: from concurrent.futures import ProcessPoolExecutor geohash_list = list(df.index) speed_factor_list = list(df.speed_factor) neighbours_list = list(df.neighbours) with tqdm(desc="生成Edge List", total=len(df.index), colour="green") as pbar: with ProcessPoolExecutor(os.cpu_count()) as executor: result = executor.map(create_edge_list, geohash_list, speed_factor_list, neighbours_list, chunksize=len(df.index)//(os.cpu_count())) for edge_list in result: elist.extend(edge_list) pbar.update(1) else: for geohash, row in tqdm(df.iterrows(), desc="生成Edge List", total=len(df.index), colour="green"): edge_list = create_edge_list(geohash, row.speed_factor, row.neighbours) elist.extend(edge_list)
优化建议
1. 预计算所有Geohash坐标,避免重复解码
当前haversine_distance函数每次计算都会重复解码Geohash,对于700万行数据来说是巨大的性能浪费。提前把所有涉及的Geohash坐标存入字典:
# 预计算所有Geohash的坐标(包含主Geohash和所有邻居) all_geohashes = set(df.index) for neighs in df.neighbours: all_geohashes.update(neighs) geo_coords = {gh_val: gh.decode(gh_val) for gh_val in all_geohashes} # 修改哈弗辛距离计算函数 def haversine_distance(geohash1, geohash2, geo_coords): if isinstance(geohash2, list): return [round(haversine_distance(geohash1, gh_val, geo_coords), 3) for gh_val in geohash2] lat1, lon1 = geo_coords[geohash1] lat2, lon2 = geo_coords[geohash2] lat1, lon1 = (float(lat1), float(lon1)) lat2, lon2 = (float(lat2), float(lon2)) lon1, lat1, lon2, lat2 = map(radians, [lon1, lat1, lon2, lat2]) dlon = lon2 - lon1 dlat = lat2 - lat1 a = sin(dlat/2)**2 + cos(lat1) * cos(lat2) * sin(dlon/2)**2 c = 2 * asin(sqrt(a)) r = 6371 return c * r
2. 用Pandas向量化计算替代Python循环
利用Pandas的向量化操作彻底摆脱Python层面的循环,性能提升显著:
# 展开邻居列表为多行数据 df_expanded = df.explode("neighbours").reset_index() # 批量获取坐标 df_expanded[["lat1", "lon1"]] = df_expanded["geohash"].map(geo_coords).apply(pd.Series) df_expanded[["lat2", "lon2"]] = df_expanded["neighbours"].map(geo_coords).apply(pd.Series) # 向量化计算哈弗辛距离 def haversine_vectorized(lat1, lon1, lat2, lon2): lat1, lon1, lat2, lon2 = map(radians, [lat1, lon1, lat2, lon2]) dlon = lon2 - lon1 dlat = lat2 - lat1 a = sin(dlat/2)**2 + cos(lat1) * cos(lat2) * sin(dlon/2)**2 c = 2 * asin(sqrt(a)) return c * 6371 df_expanded["distance"] = haversine_vectorized(df_expanded["lat1"], df_expanded["lon1"], df_expanded["lat2"], df_expanded["lon2"]).round(3) # 批量计算时间 speed_multiplier = 60 / (3.6 * 14) df_expanded["time"] = (df_expanded["distance"] / df_expanded["speed_factor"] * speed_multiplier).round(2) # 生成最终edge_list elist = list(df_expanded.apply(lambda row: (row["geohash"], row["neighbours"], {"distance": row["distance"], "time": row["time"]}), axis=1))
3. 优化多进程的chunksize参数
当前chunksize=len(df.index)//(os.cpu_count())的设置不一定最优,CPU密集型任务适合更大的chunksize以减少进程通信开销,可以尝试设置chunksize=1000或根据实际测试调整。
4. 更换更高效的Geohash库
geohash库的解码速度偏慢,可尝试pygeohash或geohash2等更高效的库,进一步降低坐标解码的耗时。
5. 减少多进程的数据拷贝开销
多进程模式下将DataFrame列转成列表会产生大量数据拷贝,可改用multiprocessing.Pool配合分块处理,直接传递DataFrame切片而非全量列表,减少内存占用和数据传输成本。
内容的提问来源于Stack Exchange,提问作者Laende

