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

Confluent-Kafka连接MSK 3.7.x KRaft集群报错:无可用协调器Broker

解决MSK KRaft集群(IAM认证+TLS)下python-confluent-kafka生产者无可用Broker问题

问题场景

我正在使用python-confluent-kafka库创建Kafka生产者,目标集群为启用IAM认证与TLS加密的MSK 3.7.x KRaft集群(集群内部及客户端与Broker间均启用TLS)。目前生产者与消费者均出现类似报错,包括:

[thrd:main]: Topic test_topic metadata information unknown
[thrd:main]: Topic test_topic partition count is zero: should refresh metadata

核心报错为无可用Broker,以下是Producer.py代码:

class MSKTokenProvider():
    def token(self):
        token, _ = MSKAuthTokenProvider.generate_auth_token('aws-region')
        return token


tp = MSKTokenProvider()

kafka_producer = Producer({
    'bootstrap.servers': 'broker-1',
    'security.protocol': 'SASL_SSL',
    'sasl.mechanism': 'OAUTHBEARER',
    'sasl.oauthbearer.config': tp,
    'ssl.ca.location': certifi.where()
})

while True:
    message = f"hello: {time.process_time()}"
    kafka_producer.poll(0)
    kafka_producer.produce(topic, message.encode('utf-8'), callback=delivery_report_func)

问题修复方案

1. 修正OAUTHBEARER认证配置

confluent-kafka库不支持直接将自定义对象传给sasl.oauthbearer.config,需要改用oauth_cb回调函数来生成IAM认证token:

def oauth_cb(config):
    token, expiry_ms = MSKAuthTokenProvider.generate_auth_token('aws-region')
    return OAuthBearerToken(token, expiry_ms // 1000)  # 转换为秒级时间戳

2. 完善Bootstrap Servers配置

确保配置的是MSK集群的全部Broker地址(至少2个),且使用IAM认证对应的9098端口,格式示例:
'bootstrap.servers': 'broker-1:9098,broker-2:9098,broker-3:9098'

3. 补充SSL安全配置

显式开启端点身份验证,避免TLS连接异常:
'ssl.endpoint.identification.algorithm': 'HTTPS'

4. 修复生产者循环逻辑

原代码缺少定期flush操作,导致消息无法正常发送且掩盖连接问题,需添加flush并调整poll逻辑。

完整修复代码

from confluent_kafka import Producer, OAuthBearerToken, KafkaError
import time
import certifi
from aws_msk_iam_auth import MSKAuthTokenProvider

def delivery_report_func(err, msg):
    if err is not None:
        print(f"消息发送失败: {err}")
    else:
        print(f"消息发送成功到 {msg.topic()} [{msg.partition()}]")

def oauth_cb(config):
    token, expiry_ms = MSKAuthTokenProvider.generate_auth_token('aws-region')
    return OAuthBearerToken(token, expiry_ms // 1000)

# 替换为你的实际配置
BOOTSTRAP_SERVERS = "broker-1:9098,broker-2:9098,broker-3:9098"
TOPIC = "test_topic"

kafka_producer = Producer({
    'bootstrap.servers': BOOTSTRAP_SERVERS,
    'security.protocol': 'SASL_SSL',
    'sasl.mechanism': 'OAUTHBEARER',
    'oauth_cb': oauth_cb,
    'ssl.ca.location': certifi.where(),
    'ssl.endpoint.identification.algorithm': 'HTTPS'
})

try:
    while True:
        message = f"hello: {time.process_time()}"
        kafka_producer.produce(TOPIC, message.encode('utf-8'), callback=delivery_report_func)
        # 定期处理回调+flush消息
        kafka_producer.poll(0.1)
        time.sleep(1)
        kafka_producer.flush()
except KeyboardInterrupt:
    print("停止生产者")
    kafka_producer.flush()

额外排查方向

  • 检查客户端所在机器/容器的安全组是否允许访问MSK集群的9098端口
  • 确认IAM角色/用户拥有kafka:DescribeCluster、kafka:WriteData等必要权限
  • 开启debug日志排查细节:添加'debug': 'broker,security'到生产者配置,查看更详细的连接报错

内容的提问来源于stack exchange,提问作者Fergus Johnson

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 21:54:52