You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 10:57:03