Golang RabbitMQ 延迟消费队列创建教程
在 Golang 中使用 RabbitMQ 创建延迟消费队列,可以使用 RabbitMQ 的插件 rabbitmq_delayed_message_exchange。该插件允许将消息发送到延迟交换机,交换机会在指定时间后将消息发送到目标队列。
步骤如下:
- 安装
rabbitmq_delayed_message_exchange插件
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
- 创建延迟交换机
args := make(amqp.Table)
args['x-delayed-type'] = 'direct'
err := ch.ExchangeDeclare('delayed-exchange', 'x-delayed-message', true, false, false, false, args)
if err != nil {
log.Fatalf('Failed to declare exchange: %v', err)
}
- 创建目标队列
_, err = ch.QueueDeclare('target-queue', true, false, false, false, nil)
if err != nil {
log.Fatalf('Failed to declare queue: %v', err)
}
- 将目标队列绑定到延迟交换机上
err = ch.QueueBind('target-queue', '', 'delayed-exchange', false, nil)
if err != nil {
log.Fatalf('Failed to bind queue to exchange: %v', err)
}
- 发送延迟消息到延迟交换机
msg := amqp.Publishing{
ContentType: 'text/plain',
Body: []byte('Hello, World!'),
Headers: amqp.Table{'x-delay': 5000},
}
err = ch.Publish('delayed-exchange', '', false, false, msg)
if err != nil {
log.Fatalf('Failed to publish message: %v', err)
}
- 消费目标队列的消息
msgs, err := ch.Consume('target-queue', '', false, false, false, false, nil)
if err != nil {
log.Fatalf('Failed to consume messages: %v', err)
}
for msg := range msgs {
log.Printf('Received message: %s', msg.Body)
msg.Ack(false)
}
在发送消息时,将延迟时间作为消息头中的 x-delay 属性进行设置,单位为毫秒。在消费消息时,需要手动调用 Ack 方法确认消息已经被处理。
原文地址: https://www.cveoy.top/t/topic/lAlZ 著作权归作者所有。请勿转载和采集!