Databricks中Scala对接带基础认证的Confluent Schema Registry报401
问题解决步骤
核心问题原因
你当前代码的初始化方式存在版本兼容问题:Confluent Schema Registry Client 6.2.1版本中,提前手动实例化RestService再传入CachedSchemaRegistryClient时,传入的认证配置不会被应用到已创建的RestService实例的HTTP客户端上,导致请求没有携带基础认证头,触发401错误。
解决方案
1. 修正初始化代码(优先尝试)
不要提前手动创建RestService,直接使用带URL参数的CachedSchemaRegistryClient构造函数,让客户端内部自动初始化RestService并加载认证配置,修正后代码如下:
val props = Map( "basic.auth.credentials.source" -> "USER_INFO", "basic.auth.user.info" -> s"$key:$secret" ).asJava // 直接传入schema registry地址,由客户端内部初始化带认证配置的RestService val schemaRegistryClient = new CachedSchemaRegistryClient(schemaRegistryUrl, 100, props) // 后续调用逻辑不变 schemaRegistryClient.getSchemaById(id)
2. 解决Databricks集群依赖冲突
Databricks运行时内置了旧版本的Confluent相关依赖,可能会优先加载并覆盖你使用的6.2.1版本逻辑,可通过以下配置调整依赖加载优先级:
- 在集群Spark配置中添加以下参数:
spark.driver.userClassPathFirst true spark.executor.userClassPathFirst true
- 打包你的Scala应用为Fat Jar时,仅包含
kafka-schema-registry-client及其依赖,排除Spark、Hadoop等集群已提供的依赖,避免类冲突。
3. 处理出站代理配置
如果你的Databricks集群配置了出站代理,Scala HTTP客户端默认会继承系统代理配置,可能导致认证信息丢失,可通过JVM参数配置Schema Registry域名不走代理:
- 在集群Spark配置中添加以下参数,替换为你实际的Schema Registry域名:
spark.driver.extraJavaOptions -Dhttp.nonProxyHosts=*.confluent.cloud|你的Schema Registry域名 spark.executor.extraJavaOptions -Dhttp.nonProxyHosts=*.confluent.cloud|你的Schema Registry域名
4. 调试验证
如果问题仍存在,可在代码中开启HTTP请求日志,确认是否携带了正确的认证头:
import org.slf4j.LoggerFactory import ch.qos.logback.classic.{Level, Logger} // 开启Confluent Rest客户端和Apache HTTP客户端的调试日志 LoggerFactory.getLogger("io.confluent.rest").asInstanceOf[Logger].setLevel(Level.DEBUG) LoggerFactory.getLogger("org.apache.http.headers").asInstanceOf[Logger].setLevel(Level.DEBUG)
运行后查看日志中请求的Authorization头是否存在且值正确。
内容的提问来源于stack exchange,提问作者alonisser
相关产品推荐
相关产品推荐

