Spark结构化流添加current_ts列值未实时更新问题咨询
为什么Spark Structured Streaming里添加的时间戳列全是同一个值?
这其实不是Spark Structured Streaming的问题,而是你没搞清楚Driver端和Executor端代码执行时机的差异导致的。
问题根源
你用的lit(System.currentTimeMillis())里,System.currentTimeMillis()是在Driver端执行的——而且只执行一次:当你提交流作业的那一刻,这个时间戳就被计算出来并固定了。之后每个微批次处理数据时,都会复用这个固定值,所以所有行的"current_ts"列都是同一个时间。
至于你说普通批处理里这个方法正常,其实是个错觉:批处理作业是单次运行的,每次启动时都会重新计算一次这个时间戳,所以整个批处理结果里的时间戳是该次作业的启动时间;但流作业是持续运行的,多个微批次会共享一开始就固定的那个值,所以差异就显出来了。
正确的解决方法
方法1:用Spark内置的current_timestamp()函数
这是最推荐的方式,Spark专门为这种场景提供了内置函数,它会在Executor端为每一行单独计算当前时间戳:
import org.apache.spark.sql.functions.current_timestamp val updatedDF = inputDF.withColumn("current_ts", current_timestamp())
如果需要毫秒级的数值类型(而不是Timestamp类型),可以转换一下:
import org.apache.spark.sql.functions.{current_timestamp, unix_timestamp} val updatedDF = inputDF.withColumn("current_ts", unix_timestamp(current_timestamp()) * 1000)
方法2:自定义无参数UDF
如果需要更灵活的时间处理(比如指定时区),可以写个无参数UDF,它同样会在Executor端每行执行:
import org.apache.spark.sql.functions.udf val getCurrentTs = udf(() => System.currentTimeMillis()) val updatedDF = inputDF.withColumn("current_ts", getCurrentTs())
总结
本质上是Driver端代码(只执行一次)和Executor端代码(每行/每个任务执行)的区别,流处理的持续运行特性放大了这个差异,只要把时间戳的计算逻辑放到Executor端执行,就能让每一行都拿到当前的时间戳啦。
内容的提问来源于stack exchange,提问作者Nats
相关产品推荐
相关产品推荐

