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

如何指定Jar包内打包证书的相对路径——Scala-Spark Streaming对接Kafka的文件不存在异常解决

解决Dataproc集群上Spark Streaming-Kafka SSL证书路径问题

这个问题的核心在于你本地开发时用的相对路径在打包成assembly jar后失效了——证书文件是嵌入在jar包内部的,而Kafka的SSL配置要求的是文件系统上的真实路径,不是jar内部的资源路径,所以Dataproc集群的JVM找不到文件,抛出NoSuchFileException。

下面给你两种可行的解决方案,你可以根据自己的场景选择:

方案一:从Jar中提取证书到临时文件(推荐,无需额外部署)

这种方法会把打包在jar里的证书资源读取出来,写入到集群节点的临时目录中,然后给Kafka参数指定这个临时文件的路径,作业结束后自动清理临时文件。

具体代码实现如下:

import java.io.{File, FileOutputStream}
import java.nio.file.Files

// 工具方法:把Jar内的资源文件提取到临时目录
def extractResourceToTemp(resourcePath: String): String = {
  // 创建临时文件,后缀为.jks,JVM退出时自动删除
  val tempFile = Files.createTempFile("kafka-ssl-", ".jks").toFile
  tempFile.deleteOnExit()

  // 通过类加载器获取Jar内的资源流
  val inputStream = getClass.getClassLoader.getResourceAsStream(resourcePath)
  val outputStream = new FileOutputStream(tempFile)

  try {
    // 把资源流写入临时文件
    inputStream.copyTo(outputStream)
  } finally {
    // 关闭流,避免资源泄漏
    inputStream.close()
    outputStream.close()
  }

  // 返回临时文件的绝对路径
  tempFile.getAbsolutePath
}

// 提取证书到临时文件(注意路径是相对于src/main/resources的,直接写cert/xxx.jks)
val truststoreTempPath = extractResourceToTemp("cert/trust.jks")
val keystoreTempPath = extractResourceToTemp("cert/keystore.jks")

// 构建Kafka参数时使用临时文件路径
val kafkaParams = Map[String, Object](
  Bootstrap_servers -> "270618767-1-1328258624:9093",
  Key_deserializer -> classOf[StringDeserializer],
  Value_deserializer -> classOf[StringDeserializer],
  Group_id -> config.Group_id,
  Auto_offset_reset -> "earliest",
  Enable_auto_commit -> (false: java.lang.Boolean),
  Security_protocol -> "SSL",
  SSL_truststore_location -> truststoreTempPath,
  SSL_truststore_password -> "****",
  SSL_keystore_location -> keystoreTempPath,
  SSL_keystore_password -> "*****",
  SSL_key_password -> "****",
  SSL_KEYSTORE_TYPE_CONFIG -> "JKS",
  SSL_TRUSTSTORE_TYPE_CONFIG -> "JKS"
)

方案二:将证书单独部署到集群所有节点

如果你的证书需要经常更新,或者不想在代码里处理临时文件,可以把证书上传到Dataproc集群的每个节点的固定目录(比如/opt/projectname/cert/),然后直接在Kafka参数里写绝对路径:

val kafkaParams = Map[String, Object](
  // ... 其他参数保持不变
  SSL_truststore_location-> "/opt/projectname/cert/trust.jks",
  SSL_keystore_location -> "/opt/projectname/cert/keystore.jks",
  // ... 其他参数保持不变
)

注意事项:

  1. 要确保集群的所有Worker节点都有这个目录和证书文件,可以用Dataproc的**初始化脚本(Init Action)**来批量将证书上传到所有节点
  2. 要保证Spark运行的用户(通常是yarn或spark)对证书文件有读取权限

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.01 03:17:42