AWS MSK公共集群SASL/SCRAM认证连接报TOPIC_AUTHORIZATION_FAILED错误
核心问题分析
错误code 29:TOPIC_AUTHORIZATION_FAILED本质是发送消息时,当前用户没有目标主题的操作权限,结合你的配置和代码,主要存在两个关键问题:
- Pulumi的ACL资源是异步部署的,代码中直接在创建ACL后发送消息,此时权限可能还未在MSK集群生效。
- 虽然开启了
auto.create.topics.enable=true,但allow.everyone.if.no.acl.found=false的设置要求用户必须具备主题创建权限,而现有ACL仅针对已存在的numtest主题配置,若主题未创建,自动创建逻辑会因权限不足失败。
具体解决步骤
1. 等待ACL完全部署后再发送消息
修改代码逻辑,利用Pulumi的Output.apply方法确保ACL资源创建完成后再执行消息发送:
from kafka import KafkaProducer from time import sleep from json import dumps import pulumi import pulumi_kafka as kafka # 初始化Producer producer = KafkaProducer( sasl_mechanism="SCRAM-SHA-512", api_version=(0,11,5), sasl_plain_password="secret-password", request_timeout_ms=500000, max_block_ms=60000, sasl_plain_username="secret-username", security_protocol="SASL_SSL", bootstrap_servers=['b-2xxxxxxxxxx.xxxxx:9198', 'b-3xxxxxxxxx.xxxxx:9198', 'b-1xxxxxxxxxx.xxxxxx:9198'], value_serializer=lambda x: dumps(x).encode('utf-8') ) # 创建主题操作ACL test_acl = kafka.Acl("test", acl_resource_name="numtest", acl_resource_type="Topic", acl_principal="User:secret-username", acl_host="*", acl_operation="All", acl_permission_type="Allow") # 定义消息发送函数 def send_messages(_): for e in range(500): data = {'number' : e} producer.send('numtest', value=data) sleep(5) # 等待ACL创建完成后触发消息发送 test_acl.id.apply(send_messages) # 保持Pulumi进程运行 pulumi.run()
2. 补充集群级主题创建权限(推荐)
若依赖auto.create.topics.enable=true自动创建主题,需添加集群级的创建权限ACL,确保自动创建逻辑能正常执行:
# 添加集群级主题创建权限 create_topic_acl = kafka.Acl("create-topic-acl", acl_resource_name="kafka-cluster", acl_resource_type="Cluster", acl_principal="User:secret-username", acl_host="*", acl_operation="Create", acl_permission_type="Allow")
3. 验证ACL生效状态
- 登录AWS控制台,进入MSK集群的权限页面,确认配置的ACL已存在,且
Principal格式(User:你的用户名)、资源名称和类型完全正确。 - 若ACL已存在仍报错,等待2-3分钟后重试,MSK的ACL同步到所有broker存在一定延迟。
4. 检查Producer配置兼容性
确认api_version参数与你的MSK集群版本匹配:MSK当前最低支持Kafka 2.2.1,若集群版本高于0.11.5,建议将api_version调整为对应版本(如(2,2,1)),避免协议不兼容导致的隐性权限问题。
内容的提问来源于stack exchange,提问作者gundocan
相关产品推荐
相关产品推荐

