Databricks中writeStream的.start()方法使用场景与存储路径问题
Databricks writeStream:.start() 方法的使用场景与常见疑问
1. .start() 方法的使用场景区分
不需要调用.start()的场景
当你使用 .table(table_name) 方法时,这个方法内部已经隐式调用了.start(),会自动启动流并将数据写入指定的Delta表(Databricks默认使用Delta格式)。这种写法适合直接将流数据写入内置或外部表的场景,代码更简洁,无需手动管理流查询对象,比如你给出的示例:
spark.readStream .format("cloudFiles") .option("cloudFiles.format", source_format) .option("cloudFiles.schemaLocation", checkpoint_directory) .load(data_source) .writeStream .option("checkpointLocation", checkpoint_directory) .option("mergeSchema", "true") .table(table_name)
需要调用.start()的场景
有两类典型场景必须显式调用.start():
- 直接写入指定路径而非表:当你需要将流数据直接存储到文件系统的某个路径(不创建表)时,必须通过
.start(path)指定输出路径,比如示例1:(myDF .writeStream .format("delta") .option("checkpointLocation", checkpointPath) .outputMode("append") .start(path) ) - 使用自定义输出逻辑:当你通过
.foreachBatch()或.foreach()实现自定义批处理逻辑(比如Upsert、自定义数据分发等)时,必须调用.start()启动流,同时还能通过返回的StreamingQuery对象(如示例2中的query)管理流的生命周期(比如awaitTermination()等待流执行完成、stop()终止流):query = (streaming_df.writeStream .foreachBatch(streaming_merge.upsert_to_delta) .outputMode("update") .option("checkpointLocation", checkpoint_directory) .trigger(availableNow=True) .start()) query.awaitTermination()
2. 调用.start()不传入path参数时的数据存储位置
分两种情况来看:
- 使用自定义输出逻辑(如.foreachBatch()):此时
.start()不需要传入path,数据的处理完全由你定义的自定义函数控制,不会自动写入任何默认路径。所有数据的存储、更新等操作都在你的自定义逻辑中实现(比如示例2中通过upsert_to_delta函数将数据合并到指定Delta表)。 - 未指定输出目标(无.table()/.path()也无自定义逻辑):直接调用
.start()会抛出错误,因为Spark需要明确知道流数据的输出位置,这种写法不合法,必须通过.path()指定路径或.table()指定表名。
内容的提问来源于stack exchange,提问作者Mohammad
相关产品推荐
相关产品推荐

