RabbitMQ Stream消费者连接异常:Bad header/Command not implemented排查
RabbitMQ Stream 客户端连接超时错误解决方案
环境配置(docker-compose)
用户初始的RabbitMQ部署配置:
services: rabbitmq: image: rabbitmq:3-management volumes: - rabbitmq_data:/var/lib/rabbitmq ports: - "5672:5672" - "15672:15672" environment: - RABBITMQ_PLUGINS=rabbitmq_stream # 原注释符号错误,YAML需用# - RABBITMQ_DEFAULT_USER=rabbitmq - RABBITMQ_DEFAULT_PASS=rabbitmq volumes: rabbitmq_data:
Stream消费者代码(Go)
package main import ( "fmt" "log" "github.com/rabbitmq/rabbitmq-stream-go-client/pkg/amqp" "github.com/rabbitmq/rabbitmq-stream-go-client/pkg/stream" ) func main() { // 创建环境 env, err := stream.NewEnvironment( stream.NewEnvironmentOptions(). SetHost("localhost"). SetPort(5672). // 此处端口错误,Stream协议默认用5552 SetUser("rabbitmq"). SetPassword("rabbitmq"), ) if err != nil { log.Fatalf("Failed to create environment: %s", err) } // 消息处理函数 handleMessage := func(consumerContext stream.ConsumerContext, message *amqp.Message) { for _, z := range message.Data { fmt.Println("Received message:", string(z)) } } // 创建消费者 _, err = env.NewConsumer( "your_stream_name", handleMessage, stream.NewConsumerOptions().SetConsumerName("my_consumer"), ) if err != nil { log.Fatalf("Failed to create consumer: %s", err) } select {} }
报错信息
客户端日志
2023/12/14 11:03:16 [warn] - Command not implemented 0 buff:0 2023/12/14 11:03:26 [error] - timeout 10000 ns - waiting Code, operation: commandPeerProperties 2023/12/14 11:03:26 [error] - Can't set the peer-properties. Check if the stream server is running/reachable 2023/12/14 11:03:26 Failed to create environment: timeout 10000 ms - waiting Code, operation: commandPeerProperties
RabbitMQ服务端日志
wss-rabbitmq-1 | 2023-12-14 16:01:07.964018+00:00 [info] <0.9.0> Time to start RabbitMQ: 30736437 us wss-rabbitmq-1 | 2023-12-14 16:02:28.713781+00:00 [info] <0.752.0> accepting AMQP connection <0.752.0> (192.168.65.1:33513 -> 192.168.240.3:5672) wss-rabbitmq-1 | 2023-12-14 16:02:28.714007+00:00 [error] <0.752.0> closing AMQP connection <0.752.0> (192.168.65.1:33513 -> 192.168.240.3:5672): {bad_header,<<0,0,0,243,0,17,0,1>>} wss-rabbitmq-1 | 2023-12-14 16:03:16.808111+00:00 [info] <0.781.0> accepting AMQP connection <0.781.0> (192.168.65.1:33844 -> 192.168.240.3:5672) wss-rabbitmq-1 | 2023-12-14 16:03:16.808803+00:00 [error] <0.781.0> closing AMQP connection <0.781.0> (192.168.65.1:33844 -> 192.168.240.3:5672): {bad_header,<<0,0,0,243,0,17,0,1>>}
已尝试操作
手动进入容器启用Stream插件:
docker exec 1d0e959696d8 rabbitmq-plugins enable rabbitmq_stream
插件启用日志显示成功,但问题依旧。
问题根源与修复步骤
1. 协议端口不匹配(核心问题)
RabbitMQ Stream协议使用独立的端口5552,而非AMQP标准的5672端口。客户端代码中错误连接到5672端口,导致协议不兼容,服务端返回bad_header错误。
修复:
- 修改客户端代码中的端口为5552:
SetPort(5552) - 在docker-compose的
ports段添加5552端口映射:ports: - "5672:5672" - "15672:15672" - "5552:5552"
2. 插件配置持久化问题
手动启用插件后,容器重启会失效。需要在docker-compose中配置环境变量,确保容器启动时自动加载Stream插件,同时修正YAML注释符号:
修复:
修改docker-compose的environment段:
environment: - RABBITMQ_PLUGINS=rabbitmq_management,rabbitmq_stream # 明确指定要加载的插件 - RABBITMQ_DEFAULT_USER=rabbitmq - RABBITMQ_DEFAULT_PASS=rabbitmq
3. 重启容器生效
执行以下命令重启RabbitMQ容器,确保配置生效:
docker-compose down docker-compose up -d
4. 验证插件状态
进入容器确认Stream插件已加载:
docker exec <container-id> rabbitmq-plugins list
看到rabbitmq_stream状态为[E*]即表示已启用。
内容的提问来源于stack exchange,提问作者Alexander Mills
相关产品推荐
相关产品推荐

