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

