apache / apache/pulsar

Unauthorized; error code: 401 when using pulsar-io-kafka connector with schema-registry that requires auth

Open
#16,181 1 comment 0 reactions 0 assignees View on GitHub
Stale
Dominant language
Java
Stars
15.3k
Forks
3.8k
Avg merge
1d 14h
Merged PRs (30d)
160

Description

Got com.google.common.util.concurrent.UncheckedExecutionException: java.lang.RuntimeException: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Unauthorized; error code: 401
when trying to use kafka-source connector with schema registry that requires basic auth

#### Steps to reproduce
Start the connector that connects to a kafka-broker, but with KafkaAvroDeserializer which goes to schema-registry. Put a basic auth to the schema-registry
use className: org.apache.pulsar.io.kafka.KafkaBytesSource
config file of pulsar-io-kafka connector:
bootstrapServers: "localhost:9092"
topic: abcV1
groupId: "group.V1.consumer-1"
valueDeserializationClass: io.confluent.kafka.serializers.KafkaAvroDeserializer
consumerConfigProperties:
client.id: "avs.sit"
security.protocol: "SASL_SSL"
sasl.mechanism: "PLAIN"
acks: "all"
client.dns.lookup: use_all_dns_ips
sasl.jaas.config: "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"USER\" password=\"password\";"
value.deserializer: io.confluent.kafka.serializers.KafkaAvroDeserializer
basic.auth.credentials.source: USER_INFO
specific.avro.reader: true
schema.registry.url: https://your-schmea-registry-url
basic.auth.user.info: USER:PASSOWRD

#### System configuration
**Pulsar version**: 2.9.2

Already digged the source code and saw that the problem might be here:
in the KafkaBytesSource.java line 110

private void initSchemaCache(Properties props) {
KafkaAvroDeserializerConfig config = new KafkaAvroDeserializerConfig(props);
List urls = config.getSchemaRegistryUrls();
int maxSchemaObject = config.getMaxSchemasPerSubject();
SchemaRegistryClient schemaRegistryClient = new CachedSchemaRegistryClient(urls, maxSchemaObject);
log.info("initializing SchemaRegistry Client, urls:{}, maxSchemasPerSubject: {}", urls, maxSchemaObject);
schemaCache = new AvroSchemaCache(schemaRegistryClient);
}

The cached schema registry is called without properties basically, and it creates a RestService without properties too,
and the default RestService never calls configure, so even thou we are passing basic.auth.user.info they never get passed to RestService so the call towards schema-registry is made without an Authorization header

Contributor guide

Open the contributing guide

Research direction

Start in KafkaBytesSource.java around initSchemaCache at line 110, then trace how CachedSchemaRegistryClient creates and configures RestService. Reproduce the Kafka source with the supplied schema-registry authentication properties and verify that the request includes authorization and no longer returns HTTP 401.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, kafka
Domain
backend, data-engineering
Issue type
Bug
Difficulty
3/5
Estimated time
1-2 days
Activity status
Stale
Clarity
Clearly specified
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.