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

使用kafka-python通过SASL_SSL连接AWS MSK集群遇连接重置错误

AWS Lambda通过SASL_SSL连接Amazon MSK时出现连接重置错误的排查

我尝试通过kafka-python库,以SASL_SSL认证方式将AWS Lambda函数连接到Amazon MSK(托管式Apache Kafka)集群,遵循官方文档操作,但在认证过程中遇到连接错误。

使用的代码

import os
import socket
from kafka import KafkaProducer
from aws_msk_iam_sasl_signer import MSKAuthTokenProvider

KAFKA_BROKERS = [
    "b-1.xxx:9098",
    "b-2.xxx:9098",
    "b-3.xxx:9098"
]                
class MSKTokenProvider():
    def token(self):
        token, _ = MSKAuthTokenProvider.generate_auth_token('us-west-2')
        return token

tp = MSKTokenProvider()
print("Generated Token:", tp.token())
print("Client:", socket.gethostname())

producer = KafkaProducer(
    bootstrap_servers=KAFKA_BROKERS,
    security_protocol='SASL_SSL',
    sasl_mechanism='OAUTHBEARER',
    sasl_oauth_token_provider=tp,
    client_id=socket.gethostname(),
    api_version=(3, 2, 0)
)

运行时错误信息

[ERROR] 2024-10-30T23:54:49.698Z e7d86aa6-cd0c-40c4-a374-42b25946d291 <BrokerConnection node_id=bootstrap-2 host=b-3.xxx:9098 <authenticating> [IPv4 ('10.29.34.48', 9098)]]: Error receiving reply from server
Traceback (most recent call last):
  File "/opt/python/kafka/conn.py", line 803, in _try_authenticate_oauth
    data = self._recv_bytes_blocking(4)
  File "/opt/python/kafka/conn.py", line 616, in _recv_bytes_blocking
    raise ConnectionError('Connection reset during recv')
ConnectionError: Connection reset during recv

环境信息

  • MSK集群Apache Kafka版本:3.2.0
  • Lambda运行时:Python 3.9

问题

认证过程中出现连接重置的原因可能是什么?是否有官方文档未提及的额外配置或排查步骤?

附相关配置截图:

  • Lambda VPC详情
  • MSK集群安全详情
  • MSK集群网络详情

排查方向与解决方案

1. 网络连通性验证

  • VPC与安全组检查:确认Lambda与MSK集群是否处于同一VPC或已通过VPC对等连接/中转网关连通;MSK集群安全组需允许Lambda所在安全组的9098端口入站流量,Lambda安全组需允许9098端口出站流量。
  • NACL规则检查:确保VPC网络访问控制列表(NACL)允许双向的9098端口流量,无拒绝规则拦截。
  • 连通性测试:在Lambda所在VPC的EC2实例上,使用telnet b-3.xxx:9098或nc -zv b-3.xxx:9098测试端口连通性,确认网络链路正常。

2. IAM权限与认证配置

  • Lambda执行角色权限:确认角色已附加包含kafka-cluster:Connect、kafka-cluster:DescribeCluster等权限的策略,且资源指定为目标MSK集群的ARN(避免使用通配符导致权限范围问题)。
  • MSK IAM认证启用状态:从安全详情截图确认集群已启用IAM访问控制,且SASL_SSL监听端口(9098)已绑定IAM认证机制。
  • Token有效性验证:在同VPC的EC2上运行代码生成Token,使用Kafka控制台客户端(如kafka-console-producer.sh)配置SASL_OAUTHBEARER认证测试连接,排查Token生成或认证逻辑问题。

3. 客户端配置细节

  • 依赖版本兼容性:确保kafka-python版本与MSK Kafka 3.2.0兼容,aws_msk_iam_sasl_signer使用最新稳定版(避免版本不匹配导致的协议兼容性问题)。
  • 区域与端点配置:确认Token生成时指定的区域(us-west-2)与MSK集群实际区域一致;可尝试显式配置sasl_oauth_token_endpoint_url参数为对应区域的MSK IAM端点(如https://kafka.us-west-2.amazonaws.com)。
  • 客户端参数补充:添加ssl_cafile参数指定AWS根证书路径(Lambda层中需包含证书文件),或启用ssl_check_hostname=True确保SSL握手正确性。

4. MSK集群状态检查

  • 监听端口配置:确认集群已启用9098端口的SASL_SSL监听,且IAM认证为该端口的唯一或有效认证机制。
  • 集群健康状态:通过AWS控制台查看MSK集群节点状态,确认无节点故障、认证服务异常等情况。

内容的提问来源于stack exchange,提问作者shady xv

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 19:04:53