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

Spark迁移S3数据到ADLS时Cosmos DB写错误日志NPE问题咨询

报错根因

你遇到的空指针异常是因为在Executor执行的算子内调用了Spark Driver端才有的API:foreachPartition内的所有逻辑都会分发到集群的Executor节点运行,而你在retry逻辑里调用的toDF()、DataFrame.write依赖Driver端的SparkSession上下文,Executor节点不存在该实例,因此触发空指针。

方案1:Executor端直接调用Cosmos Java SDK写入错误

该方案不需要把错误拉回Driver,适合错误量较大的场景,修改后的代码如下:

// 首先引入Cosmos DB Java SDK依赖,如果是sbt项目添加:
// libraryDependencies += "com.azure" % "azure-cosmos" % "4.51.0"

import com.azure.cosmos._
import com.azure.cosmos.models._
import java.util.UUID

// 重试函数去掉Spark写入逻辑,改为传入Cosmos客户端写入
def retry[T](n: Int, cosmosClient: CosmosClient, dbName: String, containerName: String, objectKey: String)(fn: => T): T = {
  util.Try {
    fn
  } match {
    case util.Success(x) => x
    case util.Failure(t: Throwable) => {
      Thread.sleep(1000)
      if (n > 1) {
        retry(n - 1, cosmosClient, dbName, containerName, objectKey)(fn)
      } else {
        // 直接用Cosmos SDK写入错误
        val container = cosmosClient.getDatabase(dbName).getContainer(containerName)
        val errorDoc = new java.util.HashMap[String, Any]()
        errorDoc.put("id", UUID.randomUUID().toString)
        errorDoc.put("Type", "Failure")
        errorDoc.put("Description", t.toString)
        errorDoc.put("s3ObjectName", objectKey)
        container.createItem(errorDoc)
        throw t
      }
    }
  }
}

// 执行逻辑调整,把Cosmos客户端也在分区内初始化
my_dataframe.rdd.foreachPartition(partition => {
  // 初始化S3客户端
  val creds = new BasicAWSCredentials(AccessKey, SecretKey)
  val clientRegion: Regions = Regions.US_EAST_1
  val s3client  = AmazonS3ClientBuilder.standard()
    .withRegion(clientRegion)
    .withCredentials(new AWSStaticCredentialsProvider(creds))
    .build()
  
  // 初始化Cosmos客户端,每个分区复用一个
  val cosmosClient = new CosmosClientBuilder()
    .endpoint(cosmosEndpoint)
    .key(cosmosKey)
    .consistencyLevel(ConsistencyLevel.EVENTUAL)
    .buildClient()
  val dbName = "你的Cosmos库名"
  val containerName = "你的Cosmos容器名"

  try {
    partition.foreach(x => {
      val objectKey = x.getString(0)
      retry(2, cosmosClient, dbName, containerName, objectKey) {
        val i = s3client.getObject(s3bucket_name, objectKey).getObjectContent
        val inputS3String = org.apache.commons.io.IOUtils.toString(i, "UTF-8")
        val filePath = s"${data_lake_file_path}/${objectKey}"
        // 建议直接用Hadoop API写ADLS,避免写本地磁盘的问题
        val hadoopConf = new org.apache.hadoop.conf.Configuration()
        val fs = org.apache.hadoop.fs.FileSystem.get(new java.net.URI(data_lake_file_path), hadoopConf)
        val outputStream = fs.create(new org.apache.hadoop.fs.Path(filePath))
        outputStream.write(inputS3String.getBytes("UTF-8"))
        outputStream.close()
      }
    })
  } finally {
    // 关闭客户端释放资源
    s3client.shutdown()
    cosmosClient.close()
  }
})

方案2:收集错误到Driver端统一写入

该方案适合错误量较小的场景,不需要在Executor侧引入Cosmos SDK,代码如下:

// 返回处理结果,Left为错误信息,Right为处理成功的对象名
val resultRDD = my_dataframe.rdd.mapPartitions(partition => {
  val creds = new BasicAWSCredentials(AccessKey, SecretKey)
  val clientRegion: Regions = Regions.US_EAST_1
  val s3client  = AmazonS3ClientBuilder.standard()
    .withRegion(clientRegion)
    .withCredentials(new AWSStaticCredentialsProvider(creds))
    .build()

  val results = partition.map(x => {
    val objectKey = x.getString(0)
    var retryCount = 2
    var success = false
    var error: Option[Throwable] = None
    while (retryCount > 0 && !success) {
      try {
        val i = s3client.getObject(s3bucket_name, objectKey).getObjectContent
        val inputS3String = org.apache.commons.io.IOUtils.toString(i, "UTF-8")
        val filePath = s"${data_lake_file_path}/${objectKey}"
        val hadoopConf = new org.apache.hadoop.conf.Configuration()
        val fs = org.apache.hadoop.fs.FileSystem.get(new java.net.URI(data_lake_file_path), hadoopConf)
        val outputStream = fs.create(new org.apache.hadoop.fs.Path(filePath))
        outputStream.write(inputS3String.getBytes("UTF-8"))
        outputStream.close()
        success = true
      } catch {
        case t: Throwable =>
          retryCount -= 1
          Thread.sleep(1000)
          if (retryCount == 0) error = Some(t)
      }
    }
    if (success) Right(objectKey) else Left((objectKey, error.get.toString))
  }).toList

  s3client.shutdown()
  results.iterator
})

// 过滤出错误记录,转成DataFrame写入Cosmos
import spark.implicits._
val errorDF = resultRDD.filter(_.isLeft).map(_.left.get).toDF("s3ObjectName", "Description")
  .withColumn("Type", lit("Failure"))
  .withColumn("id", uuid())
errorDF.write.format("cosmos.oltp").options(ExceptionCfg).mode("APPEND").save()

额外优化建议

  • 可以单独捕获S3的404异常,这类错误不需要重试,直接记录即可,减少不必要的等待时间
  • 如果需要写入的S3对象较大,建议用流写入ADLS,不要把整个对象内容加载到内存,避免OOM
  • Cosmos SDK的配置可以根据业务场景调整吞吐量参数,避免写入限流

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 06:36:02