在 Golang 中使用 RabbitMQ 创建延迟消费队列,可以使用 RabbitMQ 的插件 rabbitmq_delayed_message_exchange。该插件允许将消息发送到延迟交换机,交换机会在指定时间后将消息发送到目标队列。

步骤如下:

  1. 安装 rabbitmq_delayed_message_exchange 插件
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
  1. 创建延迟交换机
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)
}
  1. 创建目标队列
_, err = ch.QueueDeclare('target-queue', true, false, false, false, nil)
if err != nil {
    log.Fatalf('Failed to declare queue: %v', err)
}
  1. 将目标队列绑定到延迟交换机上
err = ch.QueueBind('target-queue', '', 'delayed-exchange', false, nil)
if err != nil {
    log.Fatalf('Failed to bind queue to exchange: %v', err)
}
  1. 发送延迟消息到延迟交换机
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)
}
  1. 消费目标队列的消息
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 方法确认消息已经被处理。

Golang RabbitMQ 延迟消费队列创建教程

原文地址: https://www.cveoy.top/t/topic/lAlZ 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录