SBT无法导入Apache库问题排查及Scala连接安全HTTPS Kafka集群消费者示例
解决Kafka消费者依赖与安全集群连接问题
一、先搞定SBT依赖下载失败的问题
你遇到的两个错误根源都是依赖配置出了问题:
- 版本不存在:
6.1.0-ccs是定制化的Kafka版本,Maven中央仓库没有这个包,你的内部仓库也找不到它。建议改用官方Apache Kafka的客户端版本,版本最好和你的Kafka集群版本匹配(比如集群是3.2.x就用3.2.x的客户端)。 - SSL证书错误:下载依赖时出现PKIX路径构建失败,说明你的JVM没有信任仓库的SSL证书。如果是公司内部仓库,需要把仓库的根证书导入到JVM的信任库中。
修正后的SBT构建文件
把Kafka依赖换成官方版本,若需要连接内部仓库,记得添加仓库配置:
version := "0.1" // Spark Streaming依赖(不需要的话可以直接删掉) libraryDependencies += "org.apache.spark" %% "spark-streaming" % "3.2.0" % "provided" // 官方Kafka客户端核心依赖,版本请和集群版本对齐 libraryDependencies += "org.apache.kafka" % "kafka-clients" % "3.2.0" // 可选:如果需要Scala封装的Kafka工具类,可以添加这个 // libraryDependencies += "org.apache.kafka" %% "kafka" % "3.2.0" % "provided" scalaVersion := "2.13.6" // 若使用公司内部仓库,添加仓库地址(示例) // resolvers += "Internal Repository" at "https://your-company-repo.com/content/repositories/releases/"
解决SSL证书问题的操作步骤
如果是公司内部仓库的证书问题,按以下步骤处理:
- 导出仓库的SSL根证书:用浏览器访问仓库URL,导出根证书为
.crt格式文件 - 使用
keytool将证书导入JVM信任库:
keytool -importcert -file /path/to/your/repo.crt -keystore $JAVA_HOME/lib/security/cacerts -alias internal-repo-cert
默认信任库密码是changeit,输入密码确认导入即可。
二、安全HTTPS Kafka集群的消费者代码示例
修复依赖后,你需要在代码中添加SSL相关配置才能连接安全集群,以下是完整可运行的示例:
package main.scala.kafka import java.util import java.util.Properties import org.apache.kafka.clients.consumer.KafkaConsumer import scala.jdk.CollectionConverters._ object SecureKafkaConsumer extends App { val TOPIC = "amg-dev-time" val props = new Properties() // 基础消费者配置 props.put("bootstrap.servers", "kafka-localhost.net:9093") props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer") props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer") props.put("group.id", "amlng-dev-realtime") props.put("auto.offset.reset", "latest") // 首次消费时的偏移量策略,可选配置 // SSL安全连接配置(根据你的集群要求调整) props.put("security.protocol", "SSL") // 单向认证只需要信任库配置;双向认证需要额外添加客户端证书配置(注释掉的部分) props.put("ssl.truststore.location", "/path/to/your/truststore.jks") props.put("ssl.truststore.password", "truststore-password") // props.put("ssl.keystore.location", "/path/to/your/client.keystore.jks") // props.put("ssl.keystore.password", "keystore-password") // props.put("ssl.key.password", "key-password") val consumer = new KafkaConsumer[String, String](props) consumer.subscribe(util.Collections.singletonList(TOPIC)) println("开始消费安全Kafka集群的消息...") while(true) { val records = consumer.poll(java.time.Duration.ofMillis(100)) for (record <- records.asScala) { println(s"收到消息:主题=${record.topic()}, 分区=${record.partition()}, 偏移量=${record.offset()}, 键=${record.key()}, 值=${record.value()}") } } }
代码说明
security.protocol设置为SSL开启安全连接模式- 如果集群要求双向认证(客户端需提供证书),则打开注释的
ssl.keystore相关配置 ssl.truststore用于存放信任的CA证书,确保客户端能验证集群的SSL证书合法性- 使用
scala.jdk.CollectionConverters.asScala将Java集合转为Scala可遍历的集合(适配Scala 2.13版本)
三、验证流程
- 先执行
sbt clean compile确认依赖下载成功,无编译错误 - 确认SSL证书文件路径、密码与集群要求一致
- 运行消费者程序,查看是否能正常接收集群消息
内容的提问来源于stack exchange,提问作者amarnath harish
相关产品推荐
相关产品推荐

