如何在Python Kafka生产者中设置auto.register.schemas:false?遇403权限错误
Kafka生产者问题排查与配置调整
问题描述
- 权限错误:使用Python编写的Kafka生产者(配置从AWS Secrets Manager获取)向Topic生产数据时,触发以下错误:
Error connecting to Kafka topic KafkaError{code=_VALUE_SERIALIZATION,val=-161,str="User is denied operation Write on Subject: xxxxx-value (HTTP status code 403, SR code 40301)"}
- 配置需求:需要在生产者中设置
auto.register.schemas为false
现有代码
def schemaConfig(): # Define the value schema (or import from an avro file) schema_str = """xxx xxx """ # Parse the configuration. config_parser = ConfigParser() config_parser['default'] = { 'bootstrap.servers': kafkaBootstrapServer, 'security.protocol': kafkaSecurityProtocol, 'sasl.mechanisms': kafkaSaslMechanism, 'sasl.username': kafkaClusterApiKey, 'sasl.password': kafkaClusterSecret, 'session.timeout.ms': kafkaRequestTimeout } config_parser['producer'] = { 'acks': 'all', } config_parser['avro'] = { 'url': kafkaSchemaRegistryUrl, 'basic.auth.user.info': kafkaSchemaRegisterApiKey+':'+kafkaSchemaSecret } # config_parser.read_file(args.config_file) config = dict(config_parser["default"]) config.update(config_parser["producer"]) # Define schema registry client schema_conf = dict(config_parser["avro"]) schema_registry_client = SchemaRegistryClient(schema_conf) # Create Avro Serializer avro_serializer = AvroSerializer(schema_registry_client, schema_str, user_to_dict) # Create serializing producer config.update( {"key.serializer": StringSerializer("utf_8"), "value.serializer": avro_serializer} ) producer = SerializingProducer(config) print("Inside the kafka call") return producer
问题解答
1. 权限错误原因与解决
这个403错误说明Schema Registry的用户没有对应Subject的写入权限。默认情况下AvroSerializer会尝试自动将本地Schema注册到Registry(如果不存在),但当前用户没有Write权限,因此触发报错。
解决方向:
- 联系Schema Registry管理员,给当前用户的API Key添加对应Subject(
xxxxx-value)的Write权限; - 关闭自动注册Schema功能(即设置
auto.register.schemas=false),前提是该Subject对应的Schema已提前存在于Registry中。
2. 设置auto.register.schemas为false的方法
在创建AvroSerializer时,通过关键字参数传入该配置项即可。修改代码中创建序列化器的部分:
# Create Avro Serializer avro_serializer = AvroSerializer( schema_registry_client, schema_str, user_to_dict, auto_register_schemas=False # 添加此行关闭自动注册 )
注意:开启该配置前,必须确保目标Subject对应的Schema已经存在于Schema Registry中,否则会触发Schema未找到的错误。
修改后的完整代码
def schemaConfig(): # Define the value schema (or import from an avro file) schema_str = """xxx xxx """ # Parse the configuration. config_parser = ConfigParser() config_parser['default'] = { 'bootstrap.servers': kafkaBootstrapServer, 'security.protocol': kafkaSecurityProtocol, 'sasl.mechanisms': kafkaSaslMechanism, 'sasl.username': kafkaClusterApiKey, 'sasl.password': kafkaClusterSecret, 'session.timeout.ms': kafkaRequestTimeout } config_parser['producer'] = { 'acks': 'all', } config_parser['avro'] = { 'url': kafkaSchemaRegistryUrl, 'basic.auth.user.info': kafkaSchemaRegisterApiKey+':'+kafkaSchemaSecret } # config_parser.read_file(args.config_file) config = dict(config_parser["default"]) config.update(config_parser["producer"]) # Define schema registry client schema_conf = dict(config_parser["avro"]) schema_registry_client = SchemaRegistryClient(schema_conf) # Create Avro Serializer with auto register disabled avro_serializer = AvroSerializer( schema_registry_client, schema_str, user_to_dict, auto_register_schemas=False ) # Create serializing producer config.update( {"key.serializer": StringSerializer("utf_8"), "value.serializer": avro_serializer} ) producer = SerializingProducer(config) print("Inside the kafka call") return producer
内容的提问来源于stack exchange,提问作者Pumudu Fernando
相关产品推荐
相关产品推荐

