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

使用PySpark将指定行的时间替换为最近的目标时间

PySpark DataFrame替换指定行的最近时间

实现步骤与代码

1. 准备数据与类型转换

首先确保time列为日期类型,若原始数据是字符串格式,先转换为DateType:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when, udf
from pyspark.sql.types import DateType
from datetime import datetime

# 初始化SparkSession(若已存在可跳过)
spark = SparkSession.builder.appName("ReplaceClosestTime").getOrCreate()

# 示例DataFrame(替换为你的实际DataFrame)
data = [
    (3241, "2024-01-31", False),
    (4344, "2019-09-01", True),
    (5775, "2022-02-01", False),
    (5394, "2018-06-16", True),
    (7645, "2023-03-11", False)
]
df = spark.createDataFrame(data, ["id", "time", "replace"])
# 转换time列为日期类型
df = df.withColumn("time", col("time").cast(DateType()))

2. 提取并广播基准时间列表

筛选出replace == false的行,收集它们的time值并广播,避免重复加载提升性能:

# 获取所有replace=false的时间列表
false_time_list = [row.time for row in df.filter(col("replace") == False).select("time").collect()]
# 广播列表到所有节点
broadcast_false_times = spark.sparkContext.broadcast(false_time_list)

3. 定义计算最近时间的UDF

编写自定义函数,计算目标时间与基准时间列表中最接近的日期:

def get_closest_time(target_date):
    if not target_date:
        return None
    min_days_diff = None
    closest_date = None
    # 遍历基准时间,找到差值最小的日期
    for ref_date in broadcast_false_times.value:
        days_diff = abs((target_date - ref_date).days)
        if min_days_diff is None or days_diff < min_days_diff:
            min_days_diff = days_diff
            closest_date = ref_date
    return closest_date

# 注册UDF,指定返回类型为日期类型
closest_time_udf = udf(get_closest_time, DateType())

4. 应用替换逻辑

使用when函数对replace == true的行替换时间,其余行保留原时间:

# 生成结果DataFrame
result_df = df.withColumn(
    "time",
    when(col("replace") == True, closest_time_udf(col("time"))).otherwise(col("time"))
)

# 查看处理结果
result_df.show()

处理结果示例

执行后输出如下:

+----+----------+-------+
|  id|      time|replace|
+----+----------+-------+
|3241|2024-01-31|  false|
|4344|2022-02-01|   true|
|5775|2022-02-01|  false|
|5394|2022-02-01|   true|
|7645|2023-03-11|  false|
+----+----------+-------+

补充说明

  • 若存在多个基准时间与目标时间差值相同的情况,当前UDF会取列表中第一个出现的日期,可根据需求修改逻辑(如取最晚/最早日期)。
  • 若replace == false的行数量极大,广播方式可能占用较多内存,此时可考虑用窗口函数或笛卡尔积关联后排序取最近值(需注意性能)。

内容的提问来源于stack exchange,提问作者upabove

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 23:12:38