从Databricks向Snowflake传递Timestamp查询变更数据失败求助
问题:Databricks传递时间戳到Snowflake执行变更追踪查询时报错
场景说明
在Snowflake中直接执行变更追踪查询可正常返回结果,但通过Databricks动态传递表名和时间戳时,触发时间旅行数据不可用的报错。已确认:
- Snowflake与Databricks处于同一区域,时间格式一致
- 目标表
test已启用变更追踪 - 动态生成的查询语句格式与Snowflake中正常执行的语句一致
Snowflake中正常执行的查询
SELECT * FROM test CHANGES(INFORMATION => DEFAULT) AT(TIMESTAMP => '2023-05-03 00:43:34.885 -7000')
Databricks中使用的代码
from pyspark.sql.functions import from_utc_timestamp df = spark.read.format("snowflake").options(**options).option("query","SELECT DATEADD( minute, -240, CURRENT_TIMESTAMP()) as lastruntime").load() display(df) df = df.withColumn("lastruntime", from_utc_timestamp("lastruntime", "America/Los_Angeles")) df = df.withColumn("lastruntime", date_format("lastruntime", "yyyy-MM-dd HH:mm:ss.SSS -7000")) rows = df.collect() lastruntime = str(rows[0][0]) readquery = f"SELECT * FROM {table} CHANGES(INFORMATION => DEFAULT) AT(TIMESTAMP => '%s')" %lastruntime print(readquery) df = spark.read.format("snowflake").options(**options).option("query", readquery).load()
报错信息
net.snowflake.client.jdbc.SnowflakeSQLException: Time travel data is not available for table TEST. The requested time is either beyond the allowed time travel period or before the object creation time.
问题原因
- 时区转换逻辑错误:从Snowflake获取的
CURRENT_TIMESTAMP()本身是带时区信息的时间戳,但代码中使用from_utc_timestamp将其当作UTC时间转换为America/Los_Angeles时区,这相当于执行了两次时区偏移计算,导致最终生成的时间戳与实际需要的时间偏差过大,超出了表的时间旅行保留范围。 - 硬编码时区偏移存在风险:
date_format中手动指定-7000作为时区偏移,若后续时区规则调整(如夏令时切换),会直接导致时间戳无效。
解决方案
方案1:直接在Snowflake侧生成目标时间戳字符串
将时间计算和格式化逻辑全部放在Snowflake中完成,避免Databricks侧的时区转换错误:
df = spark.read.format("snowflake").options(**options).option("query", "SELECT TO_CHAR(DATEADD(minute, -240, CURRENT_TIMESTAMP()), 'YYYY-MM-DD HH24:MI:SS.FF3 TZH:TZM') as lastruntime" ).load() lastruntime = df.collect()[0][0] # 后续拼接查询逻辑保持不变 readquery = f"SELECT * FROM {table} CHANGES(INFORMATION => DEFAULT) AT(TIMESTAMP => '{lastruntime}')" df = spark.read.format("snowflake").options(**options).option("query", readquery).load()
方案2:修正Databricks侧的时区转换逻辑
如果需要在Databricks处理时间,需确保转换逻辑正确,避免错误的UTC假设:
from pyspark.sql.functions import date_format df = spark.read.format("snowflake").options(**options).option("query","SELECT DATEADD( minute, -240, CURRENT_TIMESTAMP()) as lastruntime").load() # 直接对带时区的时间戳做格式化,用Z自动匹配时区偏移 df = df.withColumn("lastruntime", date_format("lastruntime", "yyyy-MM-dd HH:mm:ss.SSS Z")) lastruntime = df.collect()[0][0] # 拼接查询 readquery = f"SELECT * FROM {table} CHANGES(INFORMATION => DEFAULT) AT(TIMESTAMP => '{lastruntime}')" df = spark.read.format("snowflake").options(**options).option("query", readquery).load()
内容的提问来源于stack exchange,提问作者Adi
相关产品推荐
相关产品推荐

