Databricks Notebook调用random、datetime.now总返回相同值问题
问题根因
这个现象是Spark的执行逻辑和lit()函数的特性共同导致的:
lit()是Spark用于生成常量列的函数,传入的参数会在Driver端生成执行计划阶段就完成计算,整个作业的生命周期内都会复用这个固定值,不会在后续任务执行、流数据批次处理阶段重新计算。- 你定义的
create_random()是Python原生函数,直接作为lit()的传参时,只会在你定义newEventStream变量的那一刻执行一次,之后所有对这个DataFrame的操作都会复用第一次生成的结果,只有重启Notebook重新运行变量定义代码才会重新计算。
修复方案
优先使用Spark内置函数实现需求,性能远高于自定义UDF,也能避免执行逻辑不符合预期的问题:
- 生成动态随机数:直接调用Spark原生的
rand()函数,会动态生成新的随机值:
from pyspark.sql.functions import rand newEventStream = streamingEvents.withColumn("random_val", rand())
- 获取动态当前时间:调用Spark原生的
current_timestamp()函数,会在每次执行/每个流批次处理时获取最新时间,需要格式化可以搭配date_format使用:
from pyspark.sql.functions import current_timestamp, date_format # 直接生成时间戳类型的时间列 newEventStream = streamingEvents.withColumn("loadTime", current_timestamp()) # 生成指定格式的字符串类型时间列 newEventStream = streamingEvents.withColumn("loadTime", date_format(current_timestamp(), "dd/MM/yyyy HH:mm:ss"))
如果确实需要自定义复杂逻辑,需要将Python函数注册为Spark UDF,调用时不要套lit():
from pyspark.sql.functions import udf from random import random from datetime import datetime # 注册随机数UDF,返回值为双精度浮点型 @udf("double") def create_random(): return random() # 注册当前时间UDF,返回值为字符串 @udf("string") def get_current_time(): now = datetime.now() return now.strftime("%d/%m/%Y %H:%M:%S") # 直接调用UDF生成动态列 newEventStream = streamingEvents.withColumn("loadTime", get_current_time())
内容的提问来源于stack exchange,提问作者T.UK
相关产品推荐
相关产品推荐

