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

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代码修复方案

  1. 调试场景:临时查看流式数据
    模拟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();
    
  2. 正式作业:写入目标存储/系统
    若需持续监听新文件并处理,需指定输出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(); // 保持作业运行
    
  3. 仅读取现有文件:改用批处理读取
    如果不需要持续监听新文件,直接用read()替代readStream(),即可正常调用show():

    Dataset<Row> newdata = spark.read()
        .format("parquet")
        .schema(dfsample.schema())
        .load(filePath);
    newdata.show();
    

内容的提问来源于stack exchange,提问作者Artur123

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 10:30:52