【问题标题】:failed to block rabbitmq msg using msg:= rang msgs (msgs is a channel)使用 msg 阻止 rabbitmq msg 失败:= rang msgs(msgs 是一个通道)
【发布时间】:2016-09-23 15:40:16
【问题描述】:

我复制了rabbitmq go example,稍作改动进行测试。

Example URL。它工作正常

代码结构:

 func main() {
     //dial rabbit server
     //declare channel/exange/queue
     msgs, err := ch.Consume()   //typeof(msgs)=<-chan Delivery

     forever := make(chan bool)

     go func() {
         for d := range msgs {
             log.Printf("Received a message: %s", d.Body)
         }
     }()

     log.Printf(" [*] Waiting for messages. To exit press CTRL+C")
     <-forever
 }

但是如果我把一些代码放到一个函数中,比如:

func ListenRabbit() (<-chan Delivery, error) {
     //dial rabbit server
     //declare channel/exange/queue
     msgs, err := ch.Consume()   //typeof(msgs)=<-chan Delivery
     return msgs, err
}

func main(){
    msgs, _ := ListenRabbit()
    for d := range msgs {
        log.Printf("Received a message: %s", d.Body)
    }
}

无法阻止 main() 等待来自服务器的消息。它会立即退出。原始代码和更改代码之间有什么区别吗? 非常感谢!

【问题讨论】:

    标签: go rabbitmq


    【解决方案1】:

    这是垃圾收集和关闭延迟的简单错误。

    假设您的代码库与此类似,因为您省略了示例中的代码。

    package main
    
    import (
        "log"
    
        "github.com/streadway/amqp"
    )
    
    func failOnError(err error, msg string) {
        if err != nil {
            log.Fatalf("%s: %s", msg, err)
        }
    }
    
    func ListenRabbit() (<-chan amqp.Delivery, error) {
        conn, err := amqp.Dial("amqp://guest:guest@localhost:5672/")
        failOnError(err, "Failed to connect to RabbitMQ")
        defer conn.Close()
    
        ch, err := conn.Channel()
        failOnError(err, "Failed to open a channel")
        defer ch.Close()
    
        q, err := ch.QueueDeclare(
            "hello", // name
            false,   // durable
            false,   // delete when usused
            false,   // exclusive
            false,   // no-wait
            nil,     // arguments
        )
        failOnError(err, "Failed to declare a queue")
    
        msgs, err := ch.Consume(
            q.Name, // queue
            "",     // consumer
            true,   // auto-ack
            false,  // exclusive
            false,  // no-local
            false,  // no-wait
            nil,    // args
        )
        failOnError(err, "Failed to register a consumer")
    
        return msgs, err
    }
    
    func main() {
        msgs, _ := ListenRabbit()
    
        for d := range msgs {
            log.Printf("Received a message: %s", d.Body)
        }
    
        log.Printf(" [*] Waiting for messages. To exit press CTRL+C")
    }
    

    您的问题是您在 ListenRabbit 方法上初始化与 Rabbit 的连接并同时关闭它。因此,当您在通道上进行范围时,它已经关闭。

    defer conn.Close()

    defer ch.Close()

    一旦方法 ListenRabbit 退出,这些告诉 go 在 connectionchannel 上调用 Close 方法。此外,通过在该方法中初始化连接和通道,您会将所有这些对象留作垃圾收集,因为一旦方法完成,将不会留下对它们的引用。

    您需要在 main 中初始化所有这些,使其保持打开和工作,或者您可以在方法返回值上返回连接和通道,但记住在完成后处理/关闭它们。

    rabbit git repository 上的代码示例是正确的方法,但它只是设计代码的一种方法。您需要了解一些面向对象编程的基本概念、go 中的编码(引用、延迟、垃圾回收等)以及您想要做什么,这样您才能决定使用哪种设计最好。

    现在只使用示例代码就足够了。

    【讨论】:

    • 谢谢。我忘记了“延迟”操作...感谢您对设计的评论:)
    • 不用担心,只是不要初始化连接和通道然后处理它,在方法上返回它
    猜你喜欢
    • 1970-01-01
    • 2012-09-24
    • 1970-01-01
    • 2023-02-07
    • 2011-11-22
    • 2019-09-08
    • 1970-01-01
    • 2021-11-01
    • 1970-01-01
    相关资源
    最近更新 更多