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,才能读取积压的消息。
修复步骤
- 调整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 } - 补充异常场景的消息确认:在缺少header、类型错误、反序列化失败等场景调用
msg.Nack(false, false),避免消息一直滞留队列(如需重新入队可改为msg.Nack(false, true))。 - 关闭事件和错误通道: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
相关产品推荐
相关产品推荐

