Azure Databricks Java Jar运行Autoloader抛出AnalysisException问题
问题:Azure Databricks中Java流式作业调用show()报错,Scala Notebook用display()正常
在Azure Databricks共享集群上运行Java构建的Jar作业时,通过Databricks Autoloader以readStream读取Azure Blob Storage中的Parquet文件,得到Dataset<Row>后调用show()方法,抛出org.apache.spark.sql.AnalysisException,提示流式查询需用writeStream.start()。但相同逻辑的Scala代码在Databricks Notebook中调用display()可正常运行。期望获取包含id1:int、id2:int、content:binary Schema的Dataset对象。
Java代码片段
Dataset<Row> newdata = spark.readStream().format("cloudFiles") .option("cloudFiles.subscriptionId", storagesubscriptionid) .option("cloudFiles.format", "parquet") .option("cloudFiles.tenantId", sptenantid) .option("cloudFiles.clientId", spappid) .option("cloudFiles.clientSecret", spsecret) .option("cloudFiles.resourceGroup", storageresourcegroup) .option("cloudFiles.connectionString", storagesasconnectionstring) // .option("cloudFiles.useNotifications", "true") .schema(dfsample.schema()).option("cloudFiles.includeExistingFiles", "true").load(filePath); newdata.show();
报错信息
WARN SQLExecution: Error executing delta metering org.apache.spark.sql.AnalysisException: Queries with streaming sources must be executed with writeStream.start(); cloudFiles at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.throwError(UnsupportedOperationChecker.scala:447) at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.$anonfun$checkForBatch$1(UnsupportedOperationChecker.scala:38) at org.apache.spark.sql.catalyst.analysis.UnsupportedOperationChecker$.$anonfun$checkForBatch$1$adapted(UnsupportedOperationChecker.scala:36)
Scala Notebook正常代码
val df1 = spark.readStream.format("cloudFiles").option("cloudFiles.useNotifications", "true").option("cloudFiles.subscriptionId", storagesubscriptionid) .option("cloudFiles.format", "parquet") .option("cloudFiles.tenantId", sptenantid) .option("cloudFiles.clientId", spappid) .option("cloudFiles.clientSecret", spsecret) .option("cloudFiles.resourceGroup", storageresourcegroup) .option("cloudFiles.connectionString", storagesasconnectionstring) .option("cloudFiles.useNotifications", "true") .option("cloudFiles.subscriptionId", storagesubscriptionid).schema(df_schema).option("cloudFiles.includeExistingFiles", "false").load(filePath); display(df1);
解决方案
核心原因
Spark Streaming生成的Dataset是流式数据集,不支持show()这类批处理API。而Databricks Notebook的display()方法做了封装:它会自动创建临时的内存Sink和触发器,启动流式查询并展示数据,无需手动调用writeStream.start()。
Java代码修复方案
调试场景:临时查看流式数据
模拟display()的逻辑,用内存Sink启动临时流式查询,再查询视图展示数据:import org.apache.spark.sql.streaming.Trigger; // 启动流式查询到内存视图 newdata.writeStream() .format("memory") .queryName("temp_stream_view") .trigger(Trigger.Once()) // 仅处理一次现有文件(对应includeExistingFiles=true) .start(); // 查询临时视图并展示数据 spark.sql("SELECT * FROM temp_stream_view").show();正式作业:写入目标存储/系统
若需持续监听新文件并处理,需指定输出Sink(如Delta Lake)和检查点路径:import org.apache.spark.sql.streaming.Trigger; newdata.writeStream() .format("delta") .option("path", "/dbfs/path/to/delta/table") // 输出路径 .option("checkpointLocation", "/dbfs/path/to/checkpoint") // 必须设置检查点,保证故障恢复 .trigger(Trigger.ProcessingTime("30 seconds")) // 每30秒触发一次处理 .start() .awaitTermination(); // 保持作业运行仅读取现有文件:改用批处理读取
如果不需要持续监听新文件,直接用read()替代readStream(),即可正常调用show():Dataset<Row> newdata = spark.read() .format("parquet") .schema(dfsample.schema()) .load(filePath); newdata.show();
内容的提问来源于stack exchange,提问作者Artur123
相关产品推荐
相关产品推荐

