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

Golang RabbitMQ消费者仅处理启动前消息的问题求助

问题:RabbitMQ Golang消费者无法实时处理消息,仅重启后处理积压消息

启动消费者程序后,发送服务发送的消息会进入booking队列,但无法被实时处理(既不打印对应日志,也不写入数据库);仅当重启消费者程序后,之前积压在队列中的消息才会被正常处理并写入数据库。需要消费者能同时处理程序启动前和启动后发送的消息。


主程序代码

package main

import (
    "flag"
    "fmt"
    "rabbit/bookingservice/listener"
    "rabbit/bookingservice/rest"
    "rabbit/lib/configuration"
    msgqueue_amqp "rabbit/lib/msgqueue/amqp"
    "rabbit/lib/persistence/dblayer"

    "github.com/streadway/amqp"
)

func main() {
    confPath := flag.String("config", "./configuration/config.json", "path to config file")
    flag.Parse()
    config, _ := configuration.ExtractConfiguration(*confPath)

    fmt.Println("config.AMQPMessageBroker: ", config.AMQPMessageBroker)
    fmt.Println("config.DBConnection: ", config.DBConnection)

    dbhandler, err := dblayer.NewPersistenceLayer(config.Databasetype, config.DBConnection)
    if err != nil {
        panic(err)
    }
    conn, err := amqp.Dial(config.AMQPMessageBroker)
    if err != nil {
        panic(err)
    }

    eventListener, err := msgqueue_amqp.NewAMQPEventListner(conn, "events", "booking")
    if err != nil {
        panic(err)
    }
    eventEmitter, err := msgqueue_amqp.NewAMQPEventEmitter(conn, "events")
    if err != nil {
        panic(err)
    }
    processor := &listener.EventProcessor{EventListener: eventListener, Database: dbhandler}

    go processor.ProcessEvents()

    rest.ServeAPI("0.0.0.0:8282", dbhandler, eventEmitter)
}

事件处理器代码

package listener

import (
    "log"
    "rabbit/contracts"
    "rabbit/lib/msgqueue"
    "rabbit/lib/persistence"

    "go.mongodb.org/mongo-driver/bson/primitive"
)

type EventProcessor struct {
    EventListener msgqueue.EventListener
    Database      persistence.DatabaseHandler
}

func (p *EventProcessor) ProcessEvents() error {
    log.Println("Listening to events...")

    received, errors, err := p.EventListener.Listen("eventCreated")

    if err != nil {
        return err
    }
    //defer close(received)
    //defer close(errors)

    for {
        select {
        case evt := <-received:
            p.handleEvent(evt)
            // if err := p.handleEvent(evt); err != nil {
            //  log.Printf("Error processing event: %s", err)
            // }
        case err := <-errors:
            log.Printf("Received error while consuming msg: %s", err)
        }
    }
}

func (p *EventProcessor) handleEvent(event msgqueue.Event) {
    switch e := event.(type) {
    case *contracts.EventCreatedEvent:
        log.Printf("event %s created: %s", e.ID, e)
        id, err := primitive.ObjectIDFromHex(e.ID)
        id2, err2 := primitive.ObjectIDFromHex(e.LocationID)
        if err != nil || err2 != nil {
            log.Printf("Error parsing ObjectID: %v", err)
            // Handle the error as needed
        } else {
            id, _, err := p.Database.AddEvent4Booking(persistence.Event{ID: id, Location: persistence.Location{ID: id2}})
            log.Printf("id created: %s", id)
            if err != nil {
                log.Printf(`{error: Error occured while persisting event %s}`, err)
                return
            }
        }
    case *contracts.LocationCreatedEvent:
        log.Printf("location %s created: %v", e.ID, e)
        id, err := primitive.ObjectIDFromHex(e.ID)
        if err != nil {
            log.Printf("Error parsing ObjectID: %v", err)
            // Handle the error as needed
        } else {
            p.Database.AddLocation(persistence.Location{ID: id})
        }
    default:
        log.Printf("unknown event type: %t", e)
    }
}

AMQP监听实现代码

package amqp

import (
    "encoding/json"
    "fmt"
    "rabbit/contracts"
    "rabbit/lib/msgqueue"

    "github.com/streadway/amqp"
)

type amqpEventListener struct {
    connection *amqp.Connection
    queue      string
    exchange   string
    //mapper     msgqueue.EventMapper
}

func NewAMQPEventListner(conn *amqp.Connection, exchange string, queue string) (msgqueue.EventListener, error) {
    listener := &amqpEventListener{
        connection: conn,
        queue:      queue,
        exchange:   exchange,
        //mapper:     msgqueue.NewEventMapper(),
    }
    err := listener.setup()
    if err != nil {
        return nil, err
    }
    return listener, nil
}
func (a *amqpEventListener) setup() error {
    channel, err := a.connection.Channel()
    if err != nil {
        return nil
    }
    defer channel.Close()
    err = channel.ExchangeDeclare(a.exchange, "topic", true, false, false, false, nil)
    if err != nil {
        return err
    }
    _, err = channel.QueueDeclare(a.queue, true, false, false, false, nil)
    if err != nil {
        return fmt.Errorf("could not declare queue %s: %s", a.queue, err)
    }
    return nil
}

func (a *amqpEventListener) Listen(eventNames ...string) (<-chan msgqueue.Event, <-chan error, error) {
    channel, err := a.connection.Channel()
    if err != nil {
        return nil, nil, err
    }

    defer channel.Close()
    for _, eventName := range eventNames {
        //if err := channel.QueueBind(a.queue, eventName, "events", false, nil); err != nil {
        // Create binding between queue and exchange for each listened event type
        if err := channel.QueueBind(a.queue, eventName, a.exchange, false, nil); err != nil {
            return nil, nil, err
        }
    }

    msgs, err := channel.Consume(a.queue, "", false, false, false, false, nil)
    if err != nil {
        return nil, nil, err
    }
    events := make(chan msgqueue.Event, 1)
    errors := make(chan error)

    go func() {
        //defer close(events)
        //defer close(errors)
        for msg := range msgs {
            rawEventName, ok := msg.Headers["x-event-name"]
            if !ok {
                errors <- fmt.Errorf("msg did not contain x-event-name header")
                msg.Nack(false, false)
                continue
            }
            eventName, ok := rawEventName.(string)
            if !ok {
                errors <- fmt.Errorf("x-event-name header is not a string, but %t", rawEventName)
                msg.Nack(false, false)
                continue
            }

            var event msgqueue.Event
            switch eventName {
            case "eventCreated":
                event = new(contracts.EventCreatedEvent)
            case "locationCreated":
                event = &contracts.LocationCreatedEvent{}
            case "eventBooked":
                event = &contracts.EventBookedEvent{}
            default:
                errors <- fmt.Errorf("=event type %s is unknown", eventName)
            }
            err := json.Unmarshal(msg.Body, event)
            if err != nil {
                errors <- err
                continue
            }
            events <- event
            msg.Ack(false)
        }
    }()
    return events, errors, nil

}

问题根源与修复方案

核心问题:Listen方法中defer channel.Close()导致通道提前关闭

在amqpEventListener.Listen方法里,创建Channel后立即使用defer channel.Close(),但方法返回后defer会触发关闭Channel,此时goroutine里的msgs通道被关闭,无法接收新消息。只有重启程序重新创建Channel,才能读取积压的消息。

修复步骤

  1. 调整Channel的关闭时机:移除Listen方法顶部的defer channel.Close(),改为在goroutine退出时关闭Channel:
    func (a *amqpEventListener) Listen(eventNames ...string) (<-chan msgqueue.Event, <-chan error, error) {
        channel, err := a.connection.Channel()
        if err != nil {
            return nil, nil, err
        }
    
        for _, eventName := range eventNames {
            if err := channel.QueueBind(a.queue, eventName, a.exchange, false, nil); err != nil {
                channel.Close() // 绑定失败时手动关闭通道
                return nil, nil, err
            }
        }
    
        msgs, err := channel.Consume(a.queue, "", false, false, false, false, nil)
        if err != nil {
            channel.Close() // 消费失败时手动关闭通道
            return nil, nil, err
        }
        events := make(chan msgqueue.Event, 1)
        errors := make(chan error)
    
        go func() {
            defer channel.Close() // 在goroutine退出时关闭通道
            defer close(events)
            defer close(errors)
            for msg := range msgs {
                rawEventName, ok := msg.Headers["x-event-name"]
                if !ok {
                    errors <- fmt.Errorf("msg did not contain x-event-name header")
                    msg.Nack(false, false)
                    continue
                }
                eventName, ok := rawEventName.(string)
                if !ok {
                    errors <- fmt.Errorf("x-event-name header is not a string, but %t", rawEventName)
                    msg.Nack(false, false)
                    continue
                }
    
                var event msgqueue.Event
                switch eventName {
                case "eventCreated":
                    event = new(contracts.EventCreatedEvent)
                case "locationCreated":
                    event = &contracts.LocationCreatedEvent{}
                case "eventBooked":
                    event = &contracts.EventBookedEvent{}
                default:
                    errors <- fmt.Errorf("event type %s is unknown", eventName)
                    msg.Nack(false, false) // 未知事件需确认拒绝
                    continue
                }
                err := json.Unmarshal(msg.Body, event)
                if err != nil {
                    errors <- err
                    msg.Nack(false, false) // 反序列化失败需确认拒绝
                    continue
                }
                events <- event
                msg.Ack(false)
            }
        }()
        return events, errors, nil
    }
    
  2. 补充异常场景的消息确认:在缺少header、类型错误、反序列化失败等场景调用msg.Nack(false, false),避免消息一直滞留队列(如需重新入队可改为msg.Nack(false, true))。
  3. 关闭事件和错误通道:goroutine退出时关闭events和errors通道,避免资源泄漏。

额外优化点

  • 若需处理locationCreated事件,需修改ProcessEvents中的调用:received, errors, err := p.EventListener.Listen("eventCreated", "locationCreated"),否则该类型消息会被标记为未知事件。
  • 主程序中忽略了configuration.ExtractConfiguration的错误,应补充处理:
    config, err := configuration.ExtractConfiguration(*confPath)
    if err != nil {
        panic(err)
    }
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 14:40:55