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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 23:24:51