Spark中无需读取输出表检查DataFrameWriter.save()写入结果的方法
无需读取输出表验证Spark写入Kusto结果的方法
1. 异常捕获直接判断任务状态
Spark的save()是阻塞式同步操作,Kusto Spark连接器在写入失败时会直接抛出对应异常(如连接超时、权限错误、数据格式不匹配等)。通过捕获异常即可明确写入结果,无需读表验证:
val writer = df.write.format("com.microsoft.kusto.spark.synapse.datasource") .option("spark.synapse.linkedService", linkedServiceName) .option("kustoDatabase", database) .option("kustoTable", table) try { writer.mode(mode).save() // 代码执行到此处,代表Spark层面已完成所有数据提交,Kusto会负责最终持久化 println("Kusto写入任务执行成功") } catch { case e: Exception => // 捕获到异常则说明写入失败,可记录日志或触发后续告警逻辑 println(s"Kusto写入任务失败:${e.getMessage}") throw e // 可选:重新抛出异常标记Spark任务失败 }
2. 开启连接器日志监控细节
通过调整Spark日志级别,让Kusto连接器输出详细写入日志,从中直接获取批次写入状态、成功行数、Kusto确认信息等:
// 初始化SparkSession时配置日志级别 val spark = SparkSession.builder() .appName("KustoWriteJob") .config("log4j.logger.com.microsoft.kusto.spark", "DEBUG") .getOrCreate()
开启DEBUG级别后,可通过Spark UI、集群日志系统(如YARN日志)查看写入的全流程细节,无需读表就能确认结果。
3. 调用Kusto Ingestion状态API验证
Kusto提供Ingestion状态查询API,可通过写入时指定的请求ID查询数据摄入状态,适合强一致性要求的场景:
val requestId = java.util.UUID.randomUUID().toString val writer = df.write.format("com.microsoft.kusto.spark.synapse.datasource") .option("spark.synapse.linkedService", linkedServiceName) .option("kustoDatabase", database) .option("kustoTable", table) .option("requestId", requestId) // 指定唯一请求ID用于后续查询 writer.mode(mode).save() // 调用Kusto IngestionStatus API查询该requestId的摄入状态 // 可通过Azure SDK或HTTP请求实现API调用逻辑
内容的提问来源于stack exchange,提问作者Ricky_Is_Trying_Her_Best
相关产品推荐
相关产品推荐

