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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 01:07:18