Azure Databricks中从Delta表读取的Spark DataFrame会自动刷新吗?
在Azure Databricks平台用Scala开发Spark代码时遇到困惑:从Delta表加载的静态硬件信息DataFrame(dfHardware)与流式DataFrame关联后,原本以为静态DataFrame不会自动刷新,但实验发现它包含了Notebook 2启动后写入底层Delta表的新数据。
具体场景
Notebook 1(定时任务)
// grab information and put into df df.write .format("delta") .mode("overwrite") .save(basePath + s"Metadata/HardwareInfo")
Notebook 2(持续运行)
val dfHardware = spark.read .format("delta") .load(basePath+"Metadata/HardwareInfo")
val dfStreamingInput = spark.readStream .format("delta") .option("ignoreDeletes", "true") .option("ignoreChanges", "true") .load(basePath + "DeltaTable/streaming-input-location")
val dfJoined = dfStreamingInput.join(dfHardware) .where(... criteria ...) .select(... criteria ...)
// do things with dfJoined // edit: I originally didn't mention the write dfJoined .writeStreaming .format("delta") .option("checkpointLocation", "...") .option("path","...") .outputMode("append") .trigger(Trigger.ProcessingTime("30 seconds")).start()
想确认:实验结果是否正确?静态读取的DataFrame是否确实会自动刷新?如果想禁止该行为该如何操作?
实验结果无误,静态DataFrame确实会自动刷新
你的流式查询设置了Trigger.ProcessingTime("30 seconds"),意味着Spark会每隔30秒触发一次查询执行。每次触发时,Spark会重新执行整个查询计划,包括重新读取dfHardware对应的Delta表,所以每次触发都会获取到Delta表的最新数据,看起来就像静态DataFrame自动刷新了。禁止自动刷新的几种方法
缓存静态DataFrame:将
dfHardware缓存到内存(或磁盘)中,这样每次流查询触发时就不会重新读取Delta表,而是直接使用缓存的数据。代码示例:val dfHardware = spark.read .format("delta") .load(basePath+"Metadata/HardwareInfo") .cache() // 或者用persist指定存储级别注意:如果后续需要更新缓存的数据,需要手动调用
dfHardware.unpersist()后重新读取缓存。使用Delta时间旅行固定版本:读取Delta表的特定版本或时间点的数据,确保每次读取的都是固定快照。代码示例:
// 指定版本号 val dfHardware = spark.read .format("delta") .option("versionAsOf", 5) // 替换为你需要的版本号 .load(basePath+"Metadata/HardwareInfo") // 或者指定时间戳 val dfHardware = spark.read .format("delta") .option("timestampAsOf", "2024-05-20 10:00:00") // 替换为你需要的时间戳 .load(basePath+"Metadata/HardwareInfo")转为内存集合/临时视图:如果数据量不大,可以将静态数据读取到内存集合中,再转为临时视图供流查询使用。代码示例:
// 转为临时视图 val hardwareList = spark.read .format("delta") .load(basePath+"Metadata/HardwareInfo") .collectAsList() spark.createDataFrame(hardwareList).createOrReplaceTempView("hardware_static") // 流查询时使用视图 val dfJoined = dfStreamingInput.join(spark.table("hardware_static")) .where(... criteria ...) .select(... criteria ...)注意:此方法仅适用于数据量较小的场景,避免内存溢出。
内容的提问来源于stack exchange,提问作者Troy Terry

