如何用Python的Pika连接Kubernetes中的RabbitMQ集群?
问题:Pika连接Kubernetes中RabbitMQ集群的认证失败
环境与配置
已通过Azure CLI执行kubectl apply -f definition.yaml成功创建3节点RabbitMQ集群,通过LoadBalancer暴露外部IP,且集群节点可通过rabbitmqctl cluster_status确认正常运行。
definition.yaml配置
apiVersion: rabbitmq.com/v1beta1 kind: RabbitmqCluster metadata: name: test-cluster-successful namespace: successful-cluster spec: replicas: 3 --- apiVersion: v1 kind: Service metadata: name: test-cluster-service namespace: successful-cluster spec: selector: app.kubernetes.io/name: test-cluster-successful ports: - name: amqp port: 5672 targetPort: 5672 - name: management port: 15672 targetPort: 15672 protocol: TCP type: LoadBalancer
Python连接代码
username = '<TEST_USER>' password = '<TEST_PASSWORD>' usercredentials = pika.PlainCredentials(username, password) parameters = pika.ConnectionParameters('<EXTERNAL_IP_ADDRESS>', 5672, '/', usercredentials) self.connection = pika.BlockingConnection(parameters) self.channel = self.connection.channel()
问题现象
- 未传入凭证时触发
ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN错误 - 提供正确用户名密码后仍报相同认证拒绝错误
- 可ping通外部IP,但无法通过
http://<EXTERNAL_IP_ADDRESS>:15672访问RabbitMQ管理面板
完整错误信息
ERROR pika.connection:connection.py:2054 Connection closed while authenticating indicating a probable authentication error ERROR pika.adapters.utils.connection_workflow:connection_workflow.py:291 AMQPConnector - reporting failure: AMQPConnectorAMQPHandshakeError: ProbableAuthenticationError: Client was disconnected at a connection stage indicating a probable authentication error: ("ConnectionClosedByBroker: (403) 'ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN. For details see the broker logfile.'",) ERROR pika.adapters.utils.connection_workflow:connection_workflow.py:746 AMQP connection workflow failed: AMQPConnectionWorkflowFailed: 1 exceptions in all; last exception - AMQPConnectorAMQPHandshakeError: ProbableAuthenticationError: Client was disconnected at a connection stage indicating a probable authentication error: ("ConnectionClosedByBroker: (403) 'ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN. For details see the broker logfile.'",); first exception - None. ERROR pika.adapters.utils.connection_workflow:connection_workflow.py:723 AMQPConnectionWorkflow - reporting failure: AMQPConnectionWorkflowFailed: 1 exceptions in all; last exception - AMQPConnectorAMQPHandshakeError: ProbableAuthenticationError: Client was disconnected at a connection stage indicating a probable authentication error: ("ConnectionClosedByBroker: (403) 'ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN. For details see the broker logfile.'",); first exception - None ERROR pika.adapters.blocking_connection:blocking_connection.py:450 Connection workflow failed: AMQPConnectionWorkflowFailed: 1 exceptions in all; last exception - AMQPConnectorAMQPHandshakeError: ProbableAuthenticationError: Client was disconnected at a connection stage indicating a probable authentication error: ("ConnectionClosedByBroker: (403) 'ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN. For details see the broker logfile.'",); first exception - None ERROR pika.adapters.blocking_connection:blocking_connection.py:457 Error in _create_connection(). Traceback (most recent call last): File "C:\Users\USER_NAME\repos\web-dashboard-backend\.venv\Lib\site-packages\pika\adapters\blocking_connection.py", line 451, in _create_connection raise self._reap_last_connection_workflow_error(error) pika.exceptions.ProbableAuthenticationError: ConnectionClosedByBroker: (403) 'ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN. For details see the broker logfile.' ERROR pyignite.connection:connection.py:131 Failed to perform handshake, connection to node(address=127.0.0.1, port=10800) with protocol context ProtocolContext(version=(1, 7, 0), features=4) failed: [WinError 10061] No connection could be made because the target machine actively refused it Traceback (most recent call last): File "C:\Users\USER_NAME\repos\web-dashboard-backend\dashboard\notification\message_broker.py", line 96, in create_connection connection = pika.BlockingConnection(parameters) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Users\USER_NAME\repos\web-dashboard-backend\.venv\Lib\site-packages\pika\adapters\blocking_connection.py", line 360, in __init__ self._impl = self._create_connection(parameters, _impl_class) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Users\USER_NAME\repos\web-dashboard-backend\.venv\Lib\site-packages\pika\adapters\blocking_connection.py", line 451, in _create_connection raise self._reap_last_connection_workflow_error(error) pika.exceptions.ProbableAuthenticationError: ConnectionClosedByBroker: (403) 'ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN. For details see the broker logfile.' During handling of the above exception, another exception occurred: Traceback (most recent call last): File "C:\Users\USER_NAME\repos\web-dashboard-backend\.venv\Lib\site-packages\pyignite\connection\connection.py", line 241, in connect result = self._connect_version() ^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Users\USER_NAME\repos\web-dashboard-backend\.venv\Lib\site-packages\pyignite\connection\connection.py", line 271, in _connect_version self._socket.connect((self.host, self.port)) ConnectionRefusedError: [WinError 10061] No connection could be made because the target machine actively refused it ERROR root:kafka_fallback_handler.py:45 KafkaFallbackHandler: Create connection exception: Can not connect. ERROR dashboard.notification.message_broker:message_broker.py:132 Exception while creating connection with rabbitMq Traceback (most recent call last): File "C:\Users\USER_NAME\repos\web-dashboard-backend\dashboard\notification\message_broker.py", line 96, in create_connection connection = pika.BlockingConnection(parameters) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Users\USER_NAME\repos\web-dashboard-backend\.venv\Lib\site-packages\pika\adapters\blocking_connection.py", line 360, in __init__ self._impl = self._create_connection(parameters, _impl_class) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Users\USER_NAME\repos\web-dashboard-backend\.venv\Lib\site-packages\pika\adapters\blocking_connection.py", line 451, in _create_connection raise self._reap_last_connection_workflow_error(error) pika.exceptions.ProbableAuthenticationError: ConnectionClosedByBroker: (403) 'ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN. For details see the broker logfile.'
核心需求
解决Pika连接RabbitMQ集群的认证问题
解决方案
1. 验证用户凭证与权限
进入RabbitMQ Pod内部,执行命令确认用户存在且权限配置正确:
# 查看所有用户 kubectl exec -n successful-cluster <rabbitmq-pod-name> -- rabbitmqctl list_users # 查看目标用户权限 kubectl exec -n successful-cluster <rabbitmq-pod-name> -- rabbitmqctl list_user_permissions <TEST_USER>
若用户不存在,执行以下命令创建并授予权限:
kubectl exec -n successful-cluster <rabbitmq-pod-name> -- rabbitmqctl add_user <TEST_USER> <TEST_PASSWORD> kubectl exec -n successful-cluster <rabbitmq-pod-name> -- rabbitmqctl set_user_tags <TEST_USER> administrator kubectl exec -n successful-cluster <rabbitmq-pod-name> -- rabbitmqctl set_permissions -p / <TEST_USER> ".*" ".*" ".*"
2. 排查网络与端口映射
- 确认LoadBalancer外部IP分配状态:
kubectl get svc -n successful-cluster test-cluster-service
- 检查Pod内部端口监听情况:
kubectl exec -n successful-cluster <rabbitmq-pod-name> -- netstat -tulpn | grep -E "5672|15672"
- 在集群内部测试连接,排除外部网络干扰:
kubectl run -n successful-cluster test-pod --image=pika/pika --rm -it -- python -c " import pika credentials = pika.PlainCredentials('<TEST_USER>', '<TEST_PASSWORD>') parameters = pika.ConnectionParameters('test-cluster-service', 5672, '/', credentials) connection = pika.BlockingConnection(parameters) print('Connection successful!') connection.close() "
3. 查看RabbitMQ认证日志
获取详细的认证失败原因,定位问题根源:
kubectl logs -n successful-cluster <rabbitmq-pod-name> | grep -i "access_refused\|authentication"
同时确认PLAIN认证机制已启用:
kubectl exec -n successful-cluster <rabbitmq-pod-name> -- rabbitmqctl list_authentication_mechanisms
4. 优化Python连接参数
确保虚拟主机与RabbitMQ配置一致,添加超时参数避免连接异常:
parameters = pika.ConnectionParameters( host='<EXTERNAL_IP_ADDRESS>', port=5672, virtual_host='/', # 需与RabbitMQ集群配置匹配 credentials=usercredentials, heartbeat=600, blocked_connection_timeout=300 )
内容的提问来源于stack exchange,提问作者Timmons
相关产品推荐
相关产品推荐

