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

如何用PySpark DataFrame按最新时间戳判断车辆状态并输出结果

PySpark DataFrame 解决方案:按车辆取最新记录并转换状态输出

原始数据

vehicle status  created date and time
A2345   OK      2022-12-05T04:24:06.000
A2345   NOT_OK  2022-12-05T04:25:06.000
B234    NOT_OK  2022-12-05T07:25:06.000
B234    NOT_OK  2022-12-05T06:25:06.000

需求说明

针对每个唯一的vehicle,获取其created date and time最新的记录;若该记录的status为NOT_OK则输出yes,否则输出no。


实现代码

1. 导入依赖模块

from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import col, row_number, when

2. 创建测试DataFrame(若已有现成DataFrame可跳过此步)

# 初始化SparkSession
spark = SparkSession.builder.appName("VehicleStatusCheck").getOrCreate()

# 构造测试数据
data = [
    ("A2345", "OK", "2022-12-05T04:24:06.000"),
    ("A2345", "NOT_OK", "2022-12-05T04:25:06.000"),
    ("B234", "NOT_OK", "2022-12-05T07:25:06.000"),
    ("B234", "NOT_OK", "2022-12-05T06:25:06.000")
]

df = spark.createDataFrame(data, ["vehicle", "status", "created date and time"])

3. 核心处理逻辑

# 定义窗口:按车辆分组,按时间降序排序,确保最新记录排在组内首位
window_spec = Window.partitionBy("vehicle").orderBy(col("created date and time").desc())

# 筛选每个车辆的最新记录
latest_records_df = df.withColumn("row_num", row_number().over(window_spec)) \
                     .filter(col("row_num") == 1) \
                     .drop("row_num")

# 将状态转换为要求的yes/no格式
result_df = latest_records_df.withColumn("output", when(col("status") == "NOT_OK", "yes").otherwise("no"))

# 输出最终结果
result_df.select("output").show(truncate=False)

执行结果

+------+
|output|
+------+
|yes   |
|yes   |
+------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:01:39