PySpark 3+ 中如何查找与值列表各元素最接近的DataFrame行
PySpark 3+ 匹配目标列表元素最近值行实现方案
实现思路:目标列表长度最多为数千级别,采用广播小表+交叉关联+窗口函数排序取Top1的方式实现,全流程使用PySpark原生算子,性能稳定。
核心代码实现
首先导入依赖:
from pyspark.sql import functions as F from pyspark.sql.window import Window
完整逻辑代码:
# 1. 将目标列表转换为Spark小表 lst = [10, 20, 30] target_df = spark.createDataFrame([(val,) for val in lst], schema=["target_val"]) # 2. 广播小表后和原数据关联,计算x与目标值的绝对差值 joined_df = F.broadcast(target_df).crossJoin(spark_df) \ .withColumn("diff", F.abs(F.col("x") - F.col("target_val"))) # 3. 按目标值分组取差值最小的行,差值相同时默认取x更小的行(可自行调整排序规则) window_spec = Window.partitionBy("target_val").orderBy("diff", "x") result_df = joined_df.withColumn("rank", F.row_number().over(window_spec)) \ .filter(F.col("rank") == 1) \ .drop("rank", "diff", "target_val")
调用result_df.show()即可得到符合要求的输出。
注意事项
- 若需要保留所有和目标值差值相同的匹配行,把
row_number()替换为rank()即可。 - 目标列表长度在10万以内时,该方案性能足够,广播小表可以避免大量shuffle操作。
- 若原数据量超过十亿行,可以先对原数据的x列做去重预处理,匹配到最近x之后再关联回原数据,能进一步降低计算量。
内容的提问来源于stack exchange,提问作者enneppi
相关产品推荐
相关产品推荐

