PySpark 2.0.1(Python2.6.6)计算时间差返回null的解决求助
解决Spark DataFrame计算时间差返回null的问题
首先,咱们先别用自定义UDF来处理这个问题——Spark本身就提供了更可靠、更高效的内置函数来处理时间类型转换和时间差计算,完全能避免你现在遇到的null问题。
问题根源分析
你写的UDF依赖Python的mktime和strptime,但这里存在几个容易导致null的隐患:
- 未处理null值:如果
time或time2列存在null值,UDF会直接返回null(哪怕你的示例数据里没有,生产环境也很容易遇到这种情况) - 时区不匹配:
mktime会使用本地系统时区,而Spark默认时间处理逻辑可能采用UTC,导致时间解析不一致 - UDF容错性差:一旦某行时间格式出现细微偏差,整个UDF就会返回null,而Spark内置函数的容错性要高得多
正确的解决方案
我们可以分两步走:先把字符串类型的时间转成Spark原生的timestamp类型,再计算时间差(以秒为单位):
import pyspark.sql.functions as fn # 第一步:将string类型的时间列转为timestamp类型,指定匹配的格式 df_with_ts = df.withColumn("time_ts", fn.to_timestamp(df["time"], "%Y-%m-%d %H:%M:%S")) \ .withColumn("time2_ts", fn.to_timestamp(df["time2"], "%Y-%m-%d %H:%M:%S")) # 第二步:计算时间差(秒为单位),用unix_timestamp转成秒数后相减 result_df = df_with_ts.withColumn("Diff", fn.unix_timestamp("time_ts") - fn.unix_timestamp("time2_ts")) # 如果你需要更直观的时间间隔格式(比如"13 hours 50 minutes 42 seconds"),可以直接让timestamp列相减 # result_df = df_with_ts.withColumn("Diff_interval", fn.col("time_ts") - fn.col("time2_ts")) result_df.show()
运行这段代码后,你就能得到正确的时间差数值,不会再出现null的情况。
如果你非要修复自定义UDF(不推荐)
如果坚持要使用自定义UDF,必须加上null值判断和异常捕获,避免因输入异常返回null:
import pyspark.sql.types as typ import pyspark.sql.functions as fn from pyspark.sql.functions import udf from time import mktime, strptime def diffdates(t1, t2): # 先判断输入是否为null if t1 is None or t2 is None: return None try: # 解析时间并计算秒差,转换为整数返回 time1 = mktime(strptime(t1,"%Y-%m-%d %H:%M:%S")) time2 = mktime(strptime(t2, "%Y-%m-%d %H:%M:%S")) return int(time1 - time2) except Exception as e: # 捕获解析错误,避免整个任务失败 return None dt = udf(diffdates, typ.IntegerType()) Time_Diff = df.withColumn('Diff', dt(df.time, df.time2)) Time_Diff.show()
不过还是强烈推荐使用Spark内置函数,不仅代码更简洁,性能也更好,还能避免很多潜在的坑。
内容的提问来源于stack exchange,提问作者Ahmad Senousi
相关产品推荐
相关产品推荐

