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

Centrifugo后端WebSocket消息发送失败问题排查求助

问题描述

在Docker容器部署Centrifugo后,用Go后端向user频道发送消息失败,Centrifugo日志报错:

json: cannot unmarshal ""publish","params":{"channel":"u..." into Go struct field protocol.Command.method of type int32

相关配置与代码如下:

docker-compose.yml

version: '3.7'
services:
centrifugo:
    container_name: centrifugo
    image: centrifugo/centrifugo:v3
    volumes:
      - ./dev_centrifugo_config.json:/centrifugo/config.json
    command: centrifugo -c config.json
    ports:
      - 8001:8000
    ulimits:
      nofile:
        soft: 65535
        hard: 65535

dev_centrifugo_config.json

{
    "token_hmac_secret_key": "my_secret",
    "api_key": "my_api_key",
    "admin_password": "password",
    "admin_secret": "secret",
    "admin": true,
    "allowed_origins": ["*"]
}

Go发送消息代码

import (
    "context"
    "encoding/json"
    "errors"
    "fmt"
    "github.com/centrifugal/centrifuge-go"
    "github.com/gin-gonic/gin"
    "github.com/gorilla/websocket"
    "log"
    "time"
)
func sent() {
    url := "ws://localhost:8001/connection/websocket"
    apiKey := "my_api_key"

    headers := http.Header{}
    headers.Add("Authorization", "apikey "+apiKey)

    conn, _, err := websocket.DefaultDialer.Dial(url, headers)
    if err != nil {
        log.Fatal("error connecting to Centrifugo:", err)
    }
    defer conn.Close()

    message := map[string]interface{}{
        "method": "publish",
        "params": map[string]interface{}{
            "channel": "user",
            "message": "your_message",
        },
    }

    wrappedMessageJSON, err := json.Marshal(message)
    if err != nil {
        log.Fatal("error encoding wrapped message to JSON:", err)
    }

    err = conn.WriteMessage(websocket.TextMessage, wrappedMessageJSON)
    if err != nil {
        log.Fatal("error writing message to Centrifugo:", err)
    }
}
错误原因

日志核心提示:method字段需要int32类型,但你传入了字符串"publish"。Centrifugo的WebSocket客户端协议中,命令方法采用数字枚举标识,而非字符串(比如publish对应数字2)。

另外,后端服务向Centrifugo发送消息,更推荐使用HTTP API接口而非WebSocket连接——WebSocket主要面向前端客户端,HTTP API才是后端集成的标准方案。

解决方案

推荐使用第二种方式:

方式1:修正WebSocket协议的method字段

将method的值从字符串改为对应枚举数字,例如publish对应2:

message := map[string]interface{}{
    "method": 2, // publish对应的协议枚举值
    "params": map[string]interface{}{
        "channel": "user",
        "message": "your_message",
    },
}

此方式需要记忆每个命令对应的数字,维护成本高,不推荐后端使用。

方式2:使用Centrifugo HTTP API发送消息(推荐)

直接调用Centrifugo的/api/publish HTTP接口,代码更简洁且符合后端集成规范:

import (
    "bytes"
    "encoding/json"
    "net/http"
    "log"
)

func sendCentrifugoMessage() {
    apiURL := "http://localhost:8001/api/publish"
    apiKey := "my_api_key"

    // 构造请求体,注意消息内容字段为data而非message
    reqBody := map[string]interface{}{
        "channel": "user",
        "data":    "your_message",
    }
    jsonBody, err := json.Marshal(reqBody)
    if err != nil {
        log.Fatal("marshal request body failed:", err)
    }

    // 创建请求并设置API Key认证头部
    req, err := http.NewRequest("POST", apiURL, bytes.NewBuffer(jsonBody))
    if err != nil {
        log.Fatal("create request failed:", err)
    }
    req.Header.Set("Content-Type", "application/json")
    req.Header.Set("Authorization", "apikey "+apiKey)

    // 发送请求并校验响应状态
    client := &http.Client{}
    resp, err := client.Do(req)
    if err != nil {
        log.Fatal("send request failed:", err)
    }
    defer resp.Body.Close()

    if resp.StatusCode != http.StatusOK {
        log.Fatalf("request failed with status: %s", resp.Status)
    }
}

额外注意事项

  • HTTP API请求体中,消息内容字段为data,与WebSocket协议的message字段区分开
  • 确保Centrifugo配置中的api_key与代码中一致,容器端口映射(8001→8000)正常

内容的提问来源于stack exchange,提问作者Jenny F

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 06:42:44