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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 00:43:14