Fix issue ingesting Kafka when the basic auth is enabled in the Schema Registry (#7203)

This commit is contained in:
Nahuel 2022-09-04 15:59:11 +02:00 committed by GitHub
parent bba908e33f
commit 8f1ef49cdb
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23

View File

@ -362,13 +362,10 @@ def _(connection, verbose: bool = False) -> KafkaClient:
consumer_config["group.id"] = "openmetadata-consumer"
if "auto.offset.reset" not in consumer_config:
consumer_config["auto.offset.reset"] = "earliest"
for key in connection.schemaRegistryConfig:
consumer_config["schema.registry." + key] = connection.schemaRegistryConfig[
key
]
logger.debug(f"Using Kafka consumer config: {consumer_config}")
consumer_client = AvroConsumer(consumer_config)
consumer_client = AvroConsumer(
consumer_config, schema_registry=schema_registry_client
)
return KafkaClient(
admin_client=admin_client,