Spark Scala计算两时间列时长返回null的解决方法
问题产生原因
返回全null的核心原因是时间格式不匹配导致时间解析失败,具体有两个错误点:
- 生成时间字符串时的格式模板写错:你使用的模板
yyyy/MM/dd hh:mm:ss:ss末尾重复写了两次秒占位符ss,生成的字符串是2022/05/24 03:49:01:01这种不符合常规时间规范的格式,秒值被重复拼接了两次。另外模板里的hh是12小时制占位符,处理0-23点的24小时制时间会出现解析偏差,正确的24小时制小时占位符是HH。 - 调用
to_timestamp()解析字符串时没有传入匹配的格式模板:Spark的to_timestamp函数如果不传第二个格式参数,默认会按yyyy-MM-dd HH:mm:ss的标准格式解析字符串,你存储的时间串是斜杠分隔、且末尾有重复秒值的非标准格式,函数无法识别就会返回null,两个null值做算术运算结果自然全为null。
正确实现方式
优先推荐直接用原始时间戳计算,没有格式转换开销,性能最高也不会出解析错误:
- 方案1:直接基于原始毫秒时间戳计算时长
import org.apache.spark.sql.functions._ val resDf = df1.withColumn("time_diff", ($"end_ts" - $"start_ts") / 1000 / 3600)
代码逻辑:毫秒时间戳直接做差得到毫秒级间隔,除以1000转为秒级间隔,再除以3600得到小时级间隔。
- 方案2:如果必须保留字符串格式的时间列,先修正格式模板,解析时传入完全匹配的格式参数
首先修正时间列生成逻辑,去掉重复的秒占位符,把12小时制的hh换成24小时制的HH:
import org.apache.spark.sql.functions._ // 如果不需要毫秒精度就用这个格式 val df2 = df1.withColumn("event_end_ts", from_unixtime($"end_ts"/1000, "yyyy/MM/dd HH:mm:ss")) .withColumn("event_start_ts", from_unixtime($"start_ts"/1000, "yyyy/MM/dd HH:mm:ss"))
计算时间差时必须给to_timestamp传入和生成时完全一致的格式模板,保证解析正确:
val resDf = df2.withColumn( "time_diff", ( to_timestamp($"event_end_ts", "yyyy/MM/dd HH:mm:ss").cast("long") - to_timestamp($"event_start_ts", "yyyy/MM/dd HH:mm:ss").cast("long") ) / 3600 )
注意:如果需要在字符串时间里保留毫秒精度,格式模板要写成
yyyy/MM/dd HH:mm:ss.SSS,毫秒占位符是大写的SSS,不要重复写秒占位符ss,解析时也要同步使用这个格式模板。
内容的提问来源于stack exchange,提问作者Galay
相关产品推荐
相关产品推荐

