You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

AWS MSK公共集群SASL/SCRAM认证连接报TOPIC_AUTHORIZATION_FAILED错误

解决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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.19 17:55:17