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

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证书问题的操作步骤

如果是公司内部仓库的证书问题,按以下步骤处理:

  1. 导出仓库的SSL根证书:用浏览器访问仓库URL,导出根证书为.crt格式文件
  2. 使用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版本)

三、验证流程

  1. 先执行sbt clean compile确认依赖下载成功,无编译错误
  2. 确认SSL证书文件路径、密码与集群要求一致
  3. 运行消费者程序,查看是否能正常接收集群消息

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 19:42:33