如何用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
相关产品推荐
相关产品推荐

