PySpark与Databricks环境下临时表更新自动同步到DataFrame方案咨询
实现方案
Spark DataFrame 本身的惰性计算特性已经天然支持你要的效果,不需要额外改造逻辑,只需要修正你现有代码的一处笔误即可。
核心原理
Spark DataFrame 并不存储实际的计算结果,只保存数据读取、转换的链路逻辑。只有当你触发 show()、count()、collect() 这类 action 操作时,才会真正执行计算链路拉取最新数据。只要你的查询指向持续更新的 TEMP1 视图,每次调用 action 都会自动拉取TEMP1的最新状态,不需要重复给df赋值。
修正后的代码
你之前的查询语句里表名写错了,改成指向TEMP1即可:
%python df = sqlContext.sql(''' select * from TEMP1 ''')
后续你在任意单元格直接执行df.show(),拿到的都是TEMP1的最新数据,不需要重新执行上面的赋值语句。
注意事项
- 不要给
df调用cache()、persist()方法,缓存会让Spark直接返回第一次计算的旧结果。如果之前已经加过缓存,调用df.unpersist()清除即可。 - 确认你更新
TEMP1的流任务配置正常,每2分钟的更新逻辑能正确覆盖TEMP1视图的内容即可。
内容的提问来源于stack exchange,提问作者Eduardo Mendes
相关产品推荐
相关产品推荐

