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

ActiveMQ Artemis中AMQP注解实现消息重发延迟失效问题问询

关于ActiveMQ Artemis中AMQP协议修改消息注解实现延迟重发的问题

在ActiveMQ Artemis中使用AMQP协议时,x-opt-delivery-time注解可指定消息延迟投递,发送消息时该注解能正常生效,但通过Modify Message disposition帧设置该注解时,无法实现消息重发延迟。

复现代码(Go语言)

package main

import (
    "context"
    "log"
    "time"

    "github.com/Azure/go-amqp"
)

func main() {
    // create connection
    opts := &amqp.ConnOptions{
        SASLType: amqp.SASLTypePlain("admin", "admin"),
    }
    conn, err := amqp.Dial(context.TODO(), "amqp://0.0.0.0:13001", opts)
    if err != nil {
        panic(err)
    }
    // create session
    session, err := conn.NewSession(context.TODO(), nil)
    if err != nil {
        panic(err)
    }
    // create a new sender
    sender, err := session.NewSender(context.TODO(), "test.queue", nil)
    if err != nil {
        panic(err)
    }
    // create a new receiver
    receiver, err := session.NewReceiver(context.TODO(), "test.queue", nil)
    if err != nil {
        panic(err)
    }

    // send test message
    msg := amqp.NewMessage([]byte{1})
    msg.Annotations = make(amqp.Annotations)
    msg.Annotations["x-opt-delivery-time"] = time.Now().Add(time.Second * 10).UnixMilli()
    log.Printf("sending msg")
    err = sender.Send(context.TODO(), msg, nil)
    if err != nil {
        panic(err)
    }

    // receive the test message
    msg, err = receiver.Receive(context.TODO(), nil)
    if err != nil {
        panic(err)
    }
    log.Printf("received initial msg")
    annotations := msg.Annotations
    if annotations == nil {
        annotations = make(amqp.Annotations)
    }
    annotations["x-opt-delivery-time"] = time.Now().Add(time.Second * 10).UnixMilli()
    err = receiver.ModifyMessage(context.TODO(), msg, &amqp.ModifyMessageOptions{
        DeliveryFailed:    true,
        UndeliverableHere: false,
        Annotations:       annotations,
    })
    if err != nil {
        panic(err)
    }

    // receive the modified test message
    msg, err = receiver.Receive(context.TODO(), nil)
    if err != nil {
        panic(err)
    }
    log.Printf("received modified msg")
    receiver.AcceptMessage(context.TODO(), msg)
}

程序运行输出

2024/07/28 21:02:34 sending msg
2024/07/28 21:02:44 received initial msg
2024/07/28 21:02:44 received modified msg

问题解答

这并非ActiveMQ Artemis的Bug,而是x-opt-delivery-time注解的设计特性——该注解仅在消息首次发送到队列时生效,用于初始延迟投递。当通过Modify Message操作修改已存在的消息注解时,Artemis不会重新评估这个延迟时间,因此无法触发延迟重发。

实现消息定时重发的替代方案:

  • 死信队列+延迟地址:配置队列的死信地址为一个延迟地址,当消息被标记为DeliveryFailed时,会被路由到延迟地址,在指定延迟时间后重新发回原队列。需在Artemis的broker.xml中配置延迟地址和死信规则。
  • 手动重发消息:接收并拒绝消息后,重新构造一条带有x-opt-delivery-time注解的新消息发送到队列,原消息直接确认。需注意处理消息重复投递的幂等性问题。
  • 原生延迟重发机制:通过设置队列的redelivery-delay参数,为队列中所有需重发的消息统一配置延迟时间。若需单条消息自定义延迟,可结合redelivery-delay-multiplier或自定义重发逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 06:00:09