使用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
问题
认证过程中出现连接重置的原因可能是什么?是否有官方文档未提及的额外配置或排查步骤?
附相关配置截图:
排查方向与解决方案
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
相关产品推荐
相关产品推荐




