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

Kafka Producer连接Schema Registry鉴权失败问题排查

Kafka Producer连接Schema Registry出现401未授权错误的原因分析

问题场景

在Google Dataproc上运行Scala任务,将记录推送到Kafka Topic,使用Schema Registry进行schema校验与序列化。相关代码如下:

case class KafkaProducerExample(envConfig: Config, isLocal: Boolean) {

  @transient
  lazy val logger: Logger = Logger.getLogger(classOf[KafkaProducerExample])

  def test_send(): Unit= {

    val srAuth = envConfig.Kafka.schemaRegAuthInfo

    val schemaRegistryServer: String = if (isLocal) "http://localhost:8081" else
      envConfig.Kafka.registry

    val schemaUrl = s"$schemaRegistryServer/subjects/topic_name-value/versions/1"
    println(schemaUrl)
    val apiKey = srAuth.get.split(":")(0)
    val apiSecret = srAuth.get.split(":")(1)

    // Set up the request headers with Confluent Cloud credentials
    val headers_1 = Map(
      "Content-Type" -> "application/vnd.schemaregistry.v1+json"
    )

    // Send a GET request to retrieve the schema
    val response = requests.get(schemaUrl, auth = Tuple2(apiKey, apiSecret), headers = headers_1, verifySslCerts = false)

    if (response.statusCode == 200) {
      val schema = response.text()
      println("Retrieved schema:")
      println(schema)
    } else {
      println(s"Failed to retrieve schema. Status code: ${response.statusCode}")
      println(response.text())
    }

    val kafkaProps = config.loadKafkaProducerProperties(envConfig, logger, isLocal)
    // Replace with your Kafka topic
    val topic = "topic_name"

    // SomeClass is an avro generated case class
    val bre_form_record = SomeClass(id = "123",
      batch_id = "123"
    )

    // Create Kafka producer
    val producer = new KafkaProducer[String, GenericRecord](kafkaProps)

    val avroRecord = new org.apache.avro.generic.GenericData.Record(SomeClass.SCHEMA$)
    avroRecord.put("id", "123")
    avroRecord.put("batch_id", "123")
  

    val record = new ProducerRecord[String, GenericRecord](topic, "key", avroRecord)

    producer.send(record)
    producer.close()
  }
}

当前现象

手动发送GET请求能成功从Schema Registry获取响应,但Kafka Producer客户端抛出如下错误:

Exception in thread "main" org.apache.kafka.common.errors.SerializationException: Error registering Avro schema{"type":"record","name"
Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Unauthorized; error code: 401

依赖版本

"org.apache.kafka" % "kafka-clients" % "3.4.0",
"org.apache.avro" % "avro" % "1.11.0",
"io.confluent" % "kafka-avro-serializer" % "7.1.1"

错误原因

  1. 核心问题:Kafka Avro序列化器未配置Schema Registry认证信息
    手动发送GET请求时,显式传入了Schema Registry的apiKey和apiSecret,所以能成功获取schema。但kafka-avro-serializer在与Schema Registry交互(比如自动注册schema、拉取schema)时,需要从Kafka Producer的配置中读取认证参数,而你的kafkaProps里没有添加这些配置,导致请求Schema Registry时未授权,触发401错误。

  2. 错误信息截断说明
    错误信息里的schema内容被截断只是异常输出的格式问题,并非schema本身有问题,真正的根因是后面的Unauthorized; error code: 401。

解决方法

在加载Kafka Producer配置时,添加Schema Registry的认证相关参数:

// 在kafkaProps中添加以下配置
kafkaProps.put("schema.registry.url", schemaRegistryServer)
kafkaProps.put("basic.auth.credentials.source", "USER_INFO")
kafkaProps.put("basic.auth.user.info", s"$apiKey:$apiSecret")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 12:04:53