如何在PySpark中计算距最近已审批交易的时间间隔?
计算PySpark DataFrame中每条记录与最近一次已审批交易的时间间隔
给定包含交易时间timestamp和状态status的DataFrame,需新增last_approved_time(最近一次已审批交易时间)和time_since_last_approved(时间间隔,单位秒,无前置审批则填充-1)两列,以下是实现方案:
实现步骤
- 确保
timestamp列转换为TimestampType(原始为字符串时需处理) - 使用窗口函数向前填充最近的
approved状态交易时间 - 计算当前时间与最近审批时间的秒级间隔,处理无前置审批的特殊情况
完整代码
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, last, unix_timestamp, when # 初始化SparkSession(已有则可跳过) spark = SparkSession.builder.appName("LastApprovedCalculation").getOrCreate() # 模拟用户提供的原始数据 data = [ ("2024-01-12 10:00:00", "approved"), ("2024-01-12 11:30:00", "declined"), ("2024-01-12 12:45:00", "approved"), ("2024-01-12 14:20:00", "approved"), ("2024-01-12 15:30:00", "declined"), ("2024-01-12 16:45:00", "approved"), ("2024-01-12 18:00:00", "declined"), ("2024-01-12 19:15:00", "approved"), ("2024-01-12 20:30:00", "approved"), ("2024-01-12 22:00:00", "approved") ] df = spark.createDataFrame(data, ["timestamp", "status"]) # 将字符串时间转换为Timestamp类型 df = df.withColumn("timestamp", col("timestamp").cast("timestamp")) # 定义窗口:按时间升序,范围覆盖当前行及之前所有记录 window_spec = Window.orderBy("timestamp").rowsBetween(Window.unboundedPreceding, Window.currentRow) # 填充最近审批时间:仅提取approved状态的时间,自动跳过null值向前填充 df = df.withColumn( "last_approved_time", last(col("timestamp").when(col("status") == "approved"), ignorenulls=True).over(window_spec) ) # 计算秒级时间间隔,无前置审批则设为-1 df = df.withColumn( "time_since_last_approved", when( col("last_approved_time").isNull(), -1 ).otherwise( unix_timestamp(col("timestamp")) - unix_timestamp(col("last_approved_time")) ) ) # 查看结果 df.show(truncate=False)
最终结果
| timestamp | status | last_approved_time | time_since_last_approved |
|---|---|---|---|
| 2024-01-12 10:00:00 | approved | null | -1 |
| 2024-01-12 11:30:00 | declined | 2024-01-12 10:00:00 | 5400 |
| 2024-01-12 12:45:00 | approved | 2024-01-12 10:00:00 | 9900 |
| 2024-01-12 14:20:00 | approved | 2024-01-12 12:45:00 | 5700 |
| 2024-01-12 15:30:00 | declined | 2024-01-12 14:20:00 | 3600 |
| 2024-01-12 16:45:00 | approved | 2024-01-12 14:20:00 | 8700 |
| 2024-01-12 18:00:00 | declined | 2024-01-12 16:45:00 | 4500 |
| 2024-01-12 19:15:00 | approved | 2024-01-12 16:45:00 | 9000 |
| 2024-01-12 20:30:00 | approved | 2024-01-12 19:15:00 | 4500 |
| 2024-01-12 22:00:00 | approved | 2024-01-12 20:30:00 | 5400 |
内容的提问来源于stack exchange,提问作者Rakshit Rao
相关产品推荐
相关产品推荐

