【问题标题】:How to solve Go-Stomp read timeout如何解决 Go-Stomp 读取超时
【发布时间】:2017-02-16 06:30:29
【问题描述】:

尝试使用 Go-Stomp 订阅 ActiveMQ(Apollo),但出现读取超时错误。我的应用应该每天 24 小时运行以处理传入的消息。

问题

  1. 有没有办法在队列中没有更多消息的情况下保留订阅?尝试放置 ConnOpt.HeartBeat 似乎也不起作用
  2. 为什么在读取超时后,我似乎还接受了一条消息?

以下是我的步骤:

  • 我将 1000 条消息放入输入队列中进行测试
  • 运行订阅者,代码如下
  • 订阅者读完 1000 条消息 2-3 秒后,看到错误“2016/10/07 17:12:44 Subscription 1: /queue/hflc-in: ERROR message:read timeout”。
  • 再放 1000 条消息,但似乎订阅已经关闭,因此没有消息未被处理

我的代码:

  var(
   serverAddr   = flag.String("server", "10.92.10.10:61613", "STOMP server    endpoint")
   messageCount = flag.Int("count", 10, "Number of messages to send/receive")
   inputQ       = flag.String("inputq", "/queue/hflc-in", "Input queue")
)

var options []func(*stomp.Conn) error = []func(*stomp.Conn) error{
   stomp.ConnOpt.Login("userid", "userpassword"),
   stomp.ConnOpt.Host("mybroker"),
   stomp.ConnOpt.HeartBeat(360*time.Second, 360*time.Second), // I put this but seems no impact
}

func main() {
  flag.Parse()
  jobschan := make(chan bean.Request, 10)
  //my init setup
  go getInput(1, jobschan)
}

func getInput(id int, jobschan chan bean.Request) {
   conn, err := stomp.Dial("tcp", *serverAddr, options...)

   if err != nil {
      println("cannot connect to server", err.Error())
      return
   }
   fmt.Printf("Connected %v \n", id)

   sub, err := conn.Subscribe(*inputQ, stomp.AckClient)
   if err != nil {
     println("cannot subscribe to", *inputQ, err.Error())
     return
   }

   fmt.Printf("Subscribed %v \n", id)
   var messageCount int
   for {
    msg := <-sub.C
    //expectedText := fmt.Sprintf("Message #%d", i)
    if msg != nil {

        actualText := string(msg.Body)
        
        var req bean.Request
        if actualText != "SHUTDOWN" {
            messageCount = messageCount + 1
            var err2 = easyjson.Unmarshal([]byte(actualText), &req)
            if err2 != nil {
                log.Error("Unable unmarshall", zap.Error(err))
                println("message body %v", msg.Body) // what is [0/0]0x0 ?
            } else {
                fmt.Printf("Subscriber %v received message, count %v \n  ", id, messageCount)
                jobschan <- req
            }
        } else {
            logchan <- "got some issue"
        }
    }
   }
  }

错误:

2016/10/07 17:12:44 订阅 1:/queue/hflc-in:错误消息:读取超时
[E] 2016-10-07T09:12:44Z 无法解组
消息正文 %v [0/0]0x0

【问题讨论】:

    标签: go stomp apollo


    【解决方案1】:

    通过添加这些行来解决:

    在 Apollo 中,注意到队列在几秒后清空后被删除,所以在 apollo.xml 中将 auto_delete_after 设置为几个小时,例如:

    <queue id="hflc-in" dlq="dlq-in" nak_limit="3" auto_delete_after="7200"/>
    <queue id="hflc-log" dlq="dlq-log" nak_limit="3" auto_delete_after="7200"/>
    <queue id="hflc-out" dlq="dlq-out" nak_limit="3" auto_delete_after="7200"/>
    

    在Go中,注意到go-stomp在队列中找不到任何消息后会立即放弃,所以在conn选项中,添加HeartBeat Error

    var options []func(*stomp.Conn) error = []func(*stomp.Conn) error{
       //.... original configuration
       stomp.ConnOpt.HeartBeatError(360 * time.Second),
    }
    

    但是仍然对第 2 个问题感到困惑。

    【讨论】:

      【解决方案2】:

      关于#2,我遇到了同样的问题,因为我将我的进程包装在一个无限循环中。我发现在“最后一条消息”上,您从频道订阅中获得了一条空消息的响应,但包含超时错误。这是为了能够优雅地处理断开连接。这是我在我的应用中实现它的方式

      func process(subscription *stomp.Subscription) (error) {
          log.Println("Waiting for a message")
          msg := <- subscription.C
          if msg.Body != nil {
              log.Println("Message from the queue: ", string(msg.Body))
          } else {
              log.Println("message is empty")
              log.Println("error consuming more messages", msg.Err.Error())
              return msg.Err
          }
          return nil
      }
      

      【讨论】:

        猜你喜欢
        • 2012-02-18
        • 2021-08-30
        • 2016-10-31
        • 1970-01-01
        • 2018-06-09
        • 2019-04-01
        • 1970-01-01
        • 2016-08-30
        • 2013-06-28
        相关资源
        最近更新 更多