Spark中如何用变量秒数计算时间戳并在SQL查询中引用
PySpark时间计算与变量引用问题解决
原始代码与问题
相关PySpark代码
lv_seconds_back = mv_time_horizon.select(col("max(time_horizon)") * 60).show() mv_now =spark.sql("select from_unixtime(unix_timestamp()) as mv_now") local_date_time =mv_now.select(date_format('mv_now', 'HH:mm:ss').alias("local_date_time")) lv_start =local_date_time.select(col("local_date_time") - expr("INTERVAL $lv_seconds_back seconds"))
问题1
如何使用lv_seconds_back变量中的秒数计算lv_start(即从local_date_time中减去对应秒数)?尝试使用expr(interval seconds)时,该方法仅接受数值,无法识别变量。
问题2
如何在如下Spark SQL查询中引用lv_start变量?当前写法不生效:
mt_cache_fauf_r_2= spark.sql("select mt_cache_fauf_r_temp from mt_cache_fauf_r_temp where RM_ZEITPUNKT>= ${lv_start} & RM_ZEITPUNKT <= ${lv_end}")
问题1解决方案
首先修正lv_seconds_back的获取方式:原代码中.show()仅用于打印结果,会返回None,无法作为变量使用。需要先提取具体数值:
# 从DataFrame中提取秒数数值 lv_seconds_back = mv_time_horizon.select(col("max(time_horizon)") * 60).first()[0]
计算lv_start有两种可行方式:
- 字符串格式化拼接表达式
直接将变量值拼入expr的字符串中,让Spark识别具体数值:
lv_start = local_date_time.select( col("local_date_time") - expr(f"INTERVAL {lv_seconds_back} seconds") ).alias("lv_start")
- 时间戳转换计算
先将local_date_time转为时间戳,减去对应秒数后再转回时间格式,避免字符串拼接问题:
from pyspark.sql.functions import unix_timestamp, date_format lv_start = local_date_time.select( date_format( unix_timestamp(col("local_date_time"), "HH:mm:ss") - lv_seconds_back, "HH:mm:ss" ).alias("lv_start") )
问题2解决方案
Spark SQL不支持直接用${变量名}引用Python变量,可通过以下两种方式解决:
- Python字符串格式化注入变量
先从lv_start中提取具体值,再拼入SQL语句(注意时间格式要与RM_ZEITPUNKT字段匹配):
# 提取lv_start的具体值 lv_start_val = lv_start.first()[0] # 假设lv_end已通过类似方式获取 lv_end_val = lv_end.first()[0] # 用f-string格式化SQL语句 mt_cache_fauf_r_2 = spark.sql(f""" SELECT mt_cache_fauf_r_temp FROM mt_cache_fauf_r_temp WHERE RM_ZEITPUNKT >= '{lv_start_val}' AND RM_ZEITPUNKT <= '{lv_end_val}' """)
- 注册临时视图引用
将lv_start和lv_end注册为临时视图,在SQL中通过子查询引用:
# 将lv_start注册为临时视图 lv_start.createOrReplaceTempView("temp_lv_start") # 同理注册lv_end lv_end.createOrReplaceTempView("temp_lv_end") # 在SQL中通过子查询获取变量值 mt_cache_fauf_r_2 = spark.sql(""" SELECT mt_cache_fauf_r_temp FROM mt_cache_fauf_r_temp WHERE RM_ZEITPUNKT >= (SELECT lv_start FROM temp_lv_start) AND RM_ZEITPUNKT <= (SELECT lv_end FROM temp_lv_end) """)
内容的提问来源于stack exchange,提问作者bhavana s Shetty
相关产品推荐
相关产品推荐

