PySpark中lead函数使用current_timestamp()作默认值报错及取值问题
PySpark中lead函数填充当前时间戳默认值的解决方案
问题原因
PySpark的lead函数第三个参数要求传入常量值,但F.current_timestamp()返回的是Column对象,不是具体的时间值,因此会抛出Column is not iterable错误。
三种可行解决方法
方法1:使用expr()编写SQL风格表达式
直接用SQL语法在expr()中调用lead,SQL允许默认值使用函数(current_timestamp()在SQL上下文里是合法的默认参数):
import pyspark.sql.functions as F from pyspark.sql import SparkSession, Window from pyspark.sql.types import TimestampType spark = SparkSession.builder.appName("test").getOrCreate() df = spark.createDataFrame(data=[(1,1,"2023-08-01 8:40"), (1,1,"2023-08-01 8:55")], \ schema = ["sk_subs_id", "base_stat_id", "subs_action_date"]) df = df.withColumn("subs_action_date", F.col("subs_action_date").cast(TimestampType())) w = Window.partitionBy("sk_subs_id").orderBy("subs_action_date") # 使用expr编写lead表达式 df.select('*', F.expr("lead(subs_action_date, 1, current_timestamp()) over (partition by sk_subs_id order by subs_action_date)") .alias("next_action_date")).show()
方法2:预取Driver端当前时间作为常量
如果不需要Executor端的实时时间,而是代码执行时的固定时间,可以先获取Python层面的当前时间,作为常量传入lead:
import datetime import pyspark.sql.functions as F from pyspark.sql import SparkSession, Window from pyspark.sql.types import TimestampType spark = SparkSession.builder.appName("test").getOrCreate() df = spark.createDataFrame(data=[(1,1,"2023-08-01 8:40"), (1,1,"2023-08-01 8:55")], \ schema = ["sk_subs_id", "base_stat_id", "subs_action_date"]) df = df.withColumn("subs_action_date", F.col("subs_action_date").cast(TimestampType())) # 获取当前时间作为常量值 current_ts = datetime.datetime.now() w = Window.partitionBy("sk_subs_id").orderBy("subs_action_date") df.select('*', F.lead("subs_action_date", 1, current_ts).over(w).alias("next_action_date")).show()
注意:这个时间是Driver进程运行时的时间,所有分区的null值都会填充这个固定时间,而非数据处理时的实时时间。
方法3:用coalesce替换null值
先调用不带默认值的lead,再用coalesce将结果中的null替换为current_timestamp(),逻辑更直观:
import pyspark.sql.functions as F from pyspark.sql import SparkSession, Window from pyspark.sql.types import TimestampType spark = SparkSession.builder.appName("test").getOrCreate() df = spark.createDataFrame(data=[(1,1,"2023-08-01 8:40"), (1,1,"2023-08-01 8:55")], \ schema = ["sk_subs_id", "base_stat_id", "subs_action_date"]) df = df.withColumn("subs_action_date", F.col("subs_action_date").cast(TimestampType())) w = Window.partitionBy("sk_subs_id").orderBy("subs_action_date") df.select('*', F.coalesce(F.lead("subs_action_date", 1).over(w), F.current_timestamp()) .alias("next_action_date")).show()
这种方式的时间是Executor处理数据时的实时时间,适合需要动态时间的场景。
内容的提问来源于stack exchange,提问作者Diana Oryol
相关产品推荐
相关产品推荐

