PySpark:基于整数列调整时间戳列(加减时长)问题
PySpark动态计算时间偏移列的问题
问题场景
现有如下PySpark DataFrame:
df id, duration, ts_col 'abc', 3, 2023-03-01 22:00:00 ...
需要新增ts_before(ts_col减去duration小时)和ts_after(ts_col加上duration小时)两列,预期结果:
df2 id, duration, ts_col, ts_before, ts_after 'abc', 3, 2023-03-01 22:00:00, 2023-03-01 19:00:00, 2023-03-02 01:00:00
已知固定时长的写法可行:
df2 = df.withColumn("ts_before", df.ts_col - F.expr(f"INTERVAL 3 HOUR"))
但将固定值3替换为列duration时报错,需排查原因并解决。
错误原因
你使用的INTERVAL 3 HOUR是硬编码的固定时长表达式,Spark无法直接把列名解析成动态的时长值。F.expr()里的字符串是静态SQL表达式,不能直接引用DataFrame的列变量,必须用SQL语法的列引用方式来动态拼接时长值。
正确解法
可以通过两种方式实现动态时间偏移:
方法1:SQL表达式拼接列值
直接在F.expr()里用SQL语法引用duration列,拼接成动态的INTERVAL表达式:
from pyspark.sql import functions as F df2 = df.withColumn( "ts_before", F.expr("ts_col - INTERVAL duration HOUR") ).withColumn( "ts_after", F.expr("ts_col + INTERVAL duration HOUR") )
方法2:使用F.make_interval()函数
Spark 3.0+支持make_interval()函数,直接传入动态的小时数:
from pyspark.sql import functions as F df2 = df.withColumn( "ts_before", df.ts_col - F.make_interval(hours=F.col("duration")) ).withColumn( "ts_after", df.ts_col + F.make_interval(hours=F.col("duration")) )
这两种方法都能实现基于duration列的动态时间加减,避免硬编码固定值的局限。
内容的提问来源于stack exchange,提问作者Movilla
相关产品推荐
相关产品推荐

