如何指定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", // ... 其他参数保持不变 )
注意事项:
- 要确保集群的所有Worker节点都有这个目录和证书文件,可以用Dataproc的**初始化脚本(Init Action)**来批量将证书上传到所有节点
- 要保证Spark运行的用户(通常是
yarn或spark)对证书文件有读取权限
内容的提问来源于stack exchange,提问作者amarnath harish
相关产品推荐
相关产品推荐

