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
相关产品推荐
相关产品推荐

