Scala中跨DataFrame列计算:求地点到城市的最短欧氏距离
解决方案
核心逻辑
- 第一步:对DF1和DF2执行交叉连接(CrossJoin),得到所有Place与City的配对组合
- 第二步:新增列计算每对组合的欧氏距离
- 第三步:按Place字段分组,取每组距离的最小值即为对应Place到所有City的最短距离
注意:你提供的伪代码中欧氏距离公式存在符号错误,正确公式为sqrt((X2-X1)^2 + (Y2-Y1)^2),不是减号
优化提示:如果DF2的数据量远小于DF1,可以在交叉连接时对DF2使用广播优化,大幅提升运行效率,避免大规模shuffle。
代码实现
PySpark 版本
from pyspark.sql import functions as F # 先重命名DF2的坐标列,避免和DF1的列名冲突 df2_renamed = df2.select( "City", F.col("lat").alias("city_lat"), F.col("lon").alias("city_lon") ) # 执行交叉连接,DF2数据量小时加broadcast优化 cross_df = df1.crossJoin(F.broadcast(df2_renamed)) # 计算每对组合的欧氏距离 distance_df = cross_df.withColumn( "euclidean_distance", F.sqrt( (F.col("city_lat") - F.col("lat"))**2 + (F.col("city_lon") - F.col("lon"))**2 ) ) # 按Place分组取最短距离,如需同时获取最近城市可替换下方agg逻辑 result_df = distance_df.groupBy("Place").agg( F.min("euclidean_distance").alias("min_distance") ) # 如需同时返回最近的City名称,使用如下agg逻辑 # result_df = distance_df.groupBy("Place").agg( # F.min_by(F.struct("euclidean_distance", "City"), "euclidean_distance").alias("nearest") # ).select("Place", "nearest.City", "nearest.euclidean_distance")
Pandas 版本(小数据集适用)
import pandas as pd import numpy as np # 实现交叉连接 cross_df = df1.assign(key=1).merge(df2.assign(key=1), on="key", suffixes=("_place", "_city")).drop("key", axis=1) # 计算欧氏距离 cross_df["euclidean_distance"] = np.sqrt( (cross_df["lat_city"] - cross_df["lat_place"])**2 + (cross_df["lon_city"] - cross_df["lon_place"])**2 ) # 分组取最短距离 result_df = cross_df.groupby("Place")["euclidean_distance"].min().reset_index(name="min_distance")
内容的提问来源于stack exchange,提问作者EnesK
相关产品推荐
相关产品推荐

