Spark Streaming Delta写入的两个问题:新增queryName的影响及外部Hive表创建失败排查
问题一:已运行的流能否新增queryName?
完全可以,不会有任何不良后果。
queryName本质上只是给你的流查询起一个标识性的名字,主要作用是:
- 在Spark UI的Streaming页面里更容易识别这个流任务
- 方便通过
spark.streams.get(queryName)来获取流对象进行管理(比如停止、查看状态) - 在日志里区分不同的流任务
它并不会影响流的核心运行逻辑——流的续跑依赖的是检查点目录里存储的源偏移量、状态数据等,这些和queryName完全无关。哪怕你给已经运行过的流新增或修改queryName,只要检查点目录的内容没动,流依然可以正常从上次中断的位置继续处理数据。
问题二:检查点放在数据目录下导致外部Hive表创建失败
报错原因
Delta表对其存储目录有严格要求:目录下只能包含Delta表的元数据文件夹(_delta_log)和数据文件,不能有其他无关的文件或文件夹。
你把检查点目录设为{target_table_location_directory}/_checkpoint,当流启动时会先创建这个_checkpoint文件夹并写入检查点数据。之后你的upsert方法里尝试用DeltaTable.createIfNotExists创建外部表时,发现目标目录已经存在_checkpoint这个非Delta相关的文件夹,就会抛出The associated location is not empty but it's not a Delta table的异常。
解决方案
调整检查点目录位置:把检查点和Delta表数据目录彻底分开,这是Delta Streaming的标准最佳实践。比如:
.option( "checkpointLocation", f"abfss://{target_table_location_filesystem}@{datalakename}.dfs.core.windows.net/curated/schemaname/_checkpoints/{target_table_name}", )这样检查点目录和表数据目录是同级的子目录,互相不会干扰。
修复当前已污染的目录:
- 先停止当前流任务
- 把目标表目录下的
_checkpoint文件夹移动到你新设置的检查点路径下(保证流可以续跑) - 重新执行创建外部表的逻辑,此时目标目录里只有Delta表的内容,就能创建成功了
调整执行顺序(可选):如果一定要把检查点放在靠近数据的位置,也必须先确保Delta表创建完成后再启动流。不过这种方式风险较高,不推荐——因为流启动时必然会先创建检查点目录,还是可能和表创建逻辑冲突。
备注:内容来源于stack exchange,提问作者Saugat Mukherjee

