You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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是否确实会自动刷新?如果想禁止该行为该如何操作?


解答
  1. 实验结果无误,静态DataFrame确实会自动刷新
    你的流式查询设置了Trigger.ProcessingTime("30 seconds"),意味着Spark会每隔30秒触发一次查询执行。每次触发时,Spark会重新执行整个查询计划,包括重新读取dfHardware对应的Delta表,所以每次触发都会获取到Delta表的最新数据,看起来就像静态DataFrame自动刷新了。

  2. 禁止自动刷新的几种方法

    • 缓存静态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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 13:19:58