PySpark如何计算每行的最早时间戳及对应列名?
解决方案
方法一:结构体数组法(适合多列场景)
扩展性强,后续新增日期列时只需修改数组构造部分即可。
转换日期类型
先将字符串格式的日期列转为DateType,确保时间比较的准确性:from pyspark.sql import SparkSession from pyspark.sql.types import DateType from pyspark.sql.functions import array, struct, col, min, expr # 初始化SparkSession(按需执行) spark = SparkSession.builder.appName("FindEarliestDate").getOrCreate() # 转换A/B/C列为日期类型 df = df.withColumn("A", col("A").cast(DateType())) \ .withColumn("B", col("B").cast(DateType())) \ .withColumn("C", col("C").cast(DateType()))构造日期-列名结构体数组
将每个日期列的名称和对应值打包成结构体,组合成数组:df = df.withColumn("date_pairs", array( struct(col("A").alias("date"), expr("'A'").alias("col")), struct(col("B").alias("date"), expr("'B'").alias("col")), struct(col("C").alias("date"), expr("'C'").alias("col")) ))提取最早日期及对应列名
使用min函数找到数组中日期最小的结构体,再拆分出列名和日期:df = df.withColumn("earliest_pair", min(col("date_pairs")).over()) \ .withColumn("earliest_col", col("earliest_pair.col")) \ .withColumn("earliest_date", col("earliest_pair.date")) # 筛选最终结果列 final_df = df.select("ID", "earliest_col", "earliest_date") final_df.show()
方法二:条件判断法(适合少列场景)
如果日期列数量较少,直接用least函数找最小日期,再通过when分支判断对应列名:
from pyspark.sql.functions import least, when # 先转换日期类型(同方法一) df = df.withColumn("A", col("A").cast(DateType())) \ .withColumn("B", col("B").cast(DateType())) \ .withColumn("C", col("C").cast(DateType())) # 计算最早日期和对应列名 df = df.withColumn("earliest_date", least(col("A"), col("B"), col("C"))) \ .withColumn("earliest_col", when(col("A") == col("earliest_date"), "A") .when(col("B") == col("earliest_date"), "B") .when(col("C") == col("earliest_date"), "C")) final_df = df.select("ID", "earliest_col", "earliest_date") final_df.show()
运行结果
两种方法都会得到如下输出:
+-----+------------+------------+ | ID|earliest_col|earliest_date| +-----+------------+------------+ |1fe2 | B| 2020-09-12| |3gef | B| 2019-03-04| +-----+------------+------------+
内容的提问来源于stack exchange,提问作者sol
相关产品推荐
相关产品推荐

