PySpark中处理含NULL的日期列时,使用least函数获取行最小日期对应列名的异常问题
PySpark中处理含NULL的日期列时,使用least函数获取行最小日期对应列名的异常问题
我来帮你分析下这个问题哈~你遇到的情况其实是因为原来的代码逻辑没考虑到Spark对含NULL的struct的比较规则,咱们一步步来拆解:
问题出在哪?
你原来的row_min函数是把每个列转成(value: 日期, col: 列名)的struct,然后用least取最小的struct。但Spark里,当struct的第一个字段(也就是日期值)是NULL时,这个struct和其他有非NULL日期的struct比较时,会被least优先选中——这就导致像row2里明明date3有值,结果却返回了date1的列名,完全不符合咱们想要的“忽略NULL,找非NULL里最小日期”的需求。而且全NULL的行(row4)也返回了date1,而不是我们期望的NULL。
怎么解决?
核心思路是:先把那些value为NULL的struct过滤掉,只在有实际日期值的struct里找最小的,全NULL时直接返回NULL。咱们可以修改row_min函数,用数组过滤+数组取最小的方式实现:
import pyspark.sql.functions as F def row_min(*cols): # 把每个列转成 (日期值, 列名) 的struct struct_list = [F.struct(F.col(c).alias("value"), F.lit(c).alias("col")) for c in cols] # 把这些struct放到数组里,过滤掉日期值为NULL的元素 filtered_structs = F.filter(F.array(*struct_list), lambda s: s.value.isNotNull()) # 如果过滤后还有元素,就取数组里最小的struct;空数组(全NULL)就返回NULL min_struct = F.when(F.size(filtered_structs) > 0, F.array_min(filtered_structs)).otherwise(F.lit(None)) # 提取最小struct对应的列名 return min_struct.col
然后用这个函数替换原来的row_min,再运行你的代码:
df.withColumn("output", row_min('date1', 'date2', 'date3')).show()
就能得到你想要的结果:
+----------+----------+----------+------+ | date1| date2| date3|output| +----------+----------+----------+------+ |2024-01-01|2024-01-02|2024-01-03| date1| | NULL| NULL|2024-01-03| date3| | NULL|2024-01-01|2024-01-02| date2| | NULL| NULL| NULL| NULL| |2024-01-01| NULL| NULL| date1| +----------+----------+----------+------+
为啥这个方法可行?
- 过滤NULL:用
F.filter把数组里所有value为NULL的struct都去掉,只保留有实际日期的。 - 数组取最小:
F.array_min会在剩余的struct里找最小的——这里的struct比较是先比value(日期),日期小的优先;如果日期相同才会比col(列名),完全符合咱们的需求。 - 全NULL处理:用
F.when判断过滤后的数组长度,为空时直接返回NULL,解决了row4的问题。
备注:内容来源于stack exchange,提问作者BigBen
相关产品推荐
相关产品推荐

