PySpark计算Timestamp与数组型Timestamp的日期差(天数)
实现PySpark DataFrame日期数组差值计算
需求说明
现有PySpark DataFrame包含Date(时间戳字符串)和Array_Date(时间戳字符串数组)两列,需要新增Date_diff列,存储Date与Array_Date中每个元素的天数差值(结果为整数数组)。
实现方法
方法一:使用Spark内置transform函数(Spark 3.0及以上推荐)
该方法利用Spark内置优化函数实现,性能优于自定义UDF,步骤如下:
- 导入依赖函数
from pyspark.sql import functions as F
- 构造测试DataFrame并转换日期类型
data = [ ("2020-06-25 00:00:00", ["2021-12-15 00:00:00", "2022-02-05 00:00:00"]), ("2019-11-08 00:00:00", ["2018-03-01 00:00:00", "2015-07-30 00:00:00", "2022-05-18 00:00:00"]) ] df = spark.createDataFrame(data, ["Date", "Array_Date"]) # 将字符串转换为Timestamp类型,用于日期计算 df = df.withColumn("Date", F.to_timestamp("Date")) \ .withColumn("Array_Date", F.transform("Array_Date", lambda x: F.to_timestamp(x)))
- 计算日期差值数组
# 对Array_Date中的每个元素,计算与Date的天数差 df = df.withColumn("Date_diff", F.transform("Array_Date", lambda x: F.datediff(x, F.col("Date"))))
- 查看结果
df.select("Date_diff").show(truncate=False)
输出结果:
+---------------------+ |Date_diff | +---------------------+ |[538, 590] | |[645, 1562, 922] | +---------------------+
方法二:自定义UDF(兼容Spark 3.0以下版本)
如果你的Spark版本不支持transform函数,可以用自定义UDF实现:
- 导入依赖模块并定义UDF
from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, IntegerType from datetime import datetime def calculate_date_diff(date_val, array_dates): # 兼容字符串和Timestamp类型输入 base_date = date_val if isinstance(date_val, datetime) else datetime.strptime(date_val, "%Y-%m-%d %H:%M:%S") diff_list = [] for d in array_dates: current_date = d if isinstance(d, datetime) else datetime.strptime(d, "%Y-%m-%d %H:%M:%S") diff = (current_date - base_date).days diff_list.append(diff) return diff_list date_diff_udf = F.udf(calculate_date_diff, ArrayType(IntegerType()))
- 应用UDF计算差值
df = df.withColumn("Date_diff", date_diff_udf(F.col("Date"), F.col("Array_Date")))
- 查看结果,同样会得到期望的整数数组。
内容的提问来源于stack exchange,提问作者kakdim
相关产品推荐
相关产品推荐

