Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and receive scheduled and delayed messages

Dernière mise à jour :Aug 09, 2026

Cette rubrique fournit des exemples de code pour envoyer et recevoir des messages planifiés et différés à l'aide du SDK client HTTP ApsaraMQ for RocketMQ pour Go.

Messages planifiés et différés

Les messages planifiés et différés contrôlent le moment de la remise des messages aux consommateurs :

  • Message différé : remis après un délai spécifié à compter de son envoi. Par exemple, envoyez un message maintenant pour qu'il soit livré 30 secondes plus tard.

  • Message planifié : remis à un instant précis. Par exemple, envoyez un message maintenant pour qu'il soit livré demain à 14 h 00.

Via HTTP, les deux types utilisent le même mécanisme : définissez le champ StartDeliverTime sur un horodatage Unix futur en millisecondes. Le broker conserve le message jusqu'à cet horodatage, puis le remet.

Pour plus d'informations, consultez la page Messages planifiés et messages différés.

Prérequis

Avant de commencer, assurez-vous d'avoir :

Fonctionnement

Les messages planifiés et différés utilisent tous deux le champ StartDeliverTime pour contrôler le moment de la remise. Définissez ce champ sur un horodatage Unix en millisecondes correspondant à l'instant où le broker doit remettre le message.

  • Remise différée : ajoutez le délai souhaité (en millisecondes) à l'heure actuelle.

      // Deliver 10 seconds from now
      msg.StartDeliverTime = time.Now().UTC().Unix() * 1000 + 10 * 1000
  • Remise planifiée : convertissez l'heure de remise cible en horodatage Unix en millisecondes.

      // Deliver at a specific time
      msg.StartDeliverTime = targetTime.UTC().Unix() * 1000
Si StartDeliverTime est défini sur un horodatage antérieur à l'heure actuelle, le message est remis immédiatement.

Envoyer des messages planifiés et différés

L'exemple suivant envoie quatre messages différés, chacun programmé pour être livré 10 secondes après son envoi.

Remplacez les espaces réservés suivants par vos valeurs réelles :

Espace réservé Description Exemple
${HTTP_ENDPOINT} Endpoint HTTP indiqué dans la section HTTP Endpoint de la page Instance Details http://1234567890.mqrest.cn-hangzhou.aliyuncs.com
${TOPIC} Nom du topic créé dans la console ApsaraMQ for RocketMQ delayed-msg-topic
${INSTANCE_ID} ID de l'instance. Si l'instance possède un namespace, spécifiez l'ID. Sinon, définissez-le sur null ou une chaîne vide MQ_INST_1234567890_ABCDEF
package main

import (
    "fmt"
    "time"
    "strconv"
    "os"

    "github.com/aliyunmq/mq-http-go-sdk"
)

func main() {
    // HTTP endpoint from the HTTP Endpoint section of the Instance Details page in the ApsaraMQ for RocketMQ console.
    endpoint := "${HTTP_ENDPOINT}"
    // AccessKey pair from environment variables.
    accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
    secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
    // Topic created in the ApsaraMQ for RocketMQ console.
    topic := "${TOPIC}"
    // Instance ID. If the instance has a namespace, specify the ID.
    // If the instance does not have a namespace, set this to null or an empty string.
    // You can check the namespace on the Instance Details page in the ApsaraMQ for RocketMQ console.
    instanceId := "${INSTANCE_ID}"

    client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")

    mqProducer := client.GetProducer(instanceId, topic)
    // Cyclically send four messages.
    for i := 0; i < 4; i++ {
        var msg mq_http_sdk.PublishMessageRequest

            msg = mq_http_sdk.PublishMessageRequest{
            MessageBody: "hello mq!",         // The message content.
            MessageTag:  "",                  // The message tag.
            Properties:  map[string]string{}, // The message properties.
            }
        // The message key.
            msg.MessageKey = "MessageKey"
            // The custom properties of the message.
            msg.Properties["a"] = strconv.Itoa(i)
            // Deliver 10 seconds from now. Set to a Unix timestamp in milliseconds.
            // For scheduled delivery, set this to the time difference between the scheduled point and the current time.
            msg.StartDeliverTime = time.Now().UTC().Unix() * 1000 + 10 * 1000

        ret, err := mqProducer.PublishMessage(msg)

        if err != nil {
            fmt.Println(err)
            return
        } else {
            fmt.Printf("Publish ---->\n\tMessageId:%s, BodyMD5:%s, \n", ret.MessageId, ret.MessageBodyMD5)
        }
        time.Sleep(time.Duration(100) * time.Millisecond)
    }
}

Recevoir des messages planifiés et différés

L'exemple suivant consomme des messages planifiés et différés à l'aide du sondage long (long polling). Le consommateur attend jusqu'à 3 secondes pour recevoir de nouveaux messages et traite jusqu'à 3 messages par lot.

Remplacez les espaces réservés suivants par vos valeurs réelles :

Espace réservé Description Exemple
${HTTP_ENDPOINT} Endpoint HTTP indiqué dans la section HTTP Endpoint de la page Instance Details http://1234567890.mqrest.cn-hangzhou.aliyuncs.com
${TOPIC} Nom du topic créé dans la console ApsaraMQ for RocketMQ delayed-msg-topic
${INSTANCE_ID} ID de l'instance. Si l'instance possède un namespace, spécifiez l'ID. Sinon, définissez-le sur null ou une chaîne vide MQ_INST_1234567890_ABCDEF
${GROUP_ID} ID du groupe de consommateurs créé dans la console ApsaraMQ for RocketMQ GID_delayed_consumer
Chaque topic ne prend en charge qu'un seul type de message. Un topic utilisé pour des messages planifiés ou différés ne peut pas envoyer ou recevoir d'autres types de messages.
package main

import (
    "fmt"
    "github.com/gogap/errors"
    "strings"
    "time"
    "os"

    "github.com/aliyunmq/mq-http-go-sdk"
)

func main() {
    // HTTP endpoint from the HTTP Endpoint section of the Instance Details page in the ApsaraMQ for RocketMQ console.
    endpoint := "${HTTP_ENDPOINT}"
    // AccessKey pair from environment variables.
    accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
    secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
    // Topic created in the ApsaraMQ for RocketMQ console.
    // Each topic can only send and receive messages of a specific type.
    topic := "${TOPIC}"
    // Instance ID. If the instance has a namespace, specify the ID.
    // If the instance does not have a namespace, set this to null or an empty string.
    // You can check the namespace on the Instance Details page in the ApsaraMQ for RocketMQ console.
    instanceId := "${INSTANCE_ID}"
    // Consumer group ID created in the ApsaraMQ for RocketMQ console.
    groupId := "${GROUP_ID}"

    client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")

    mqConsumer := client.GetConsumer(instanceId, topic, groupId, "")

    for {
        endChan := make(chan int)
        respChan := make(chan mq_http_sdk.ConsumeMessageResponse)
        errChan := make(chan error)
        go func() {
            select {
            case resp := <-respChan:
                {
                    // The message consumption logic.
                    var handles []string
                    fmt.Printf("Consume %d messages---->\n", len(resp.Messages))
                    for _, v := range resp.Messages {
                        handles = append(handles, v.ReceiptHandle)
                        fmt.Printf("\tMessageID: %s, PublishTime: %d, MessageTag: %s\n"+
                            "\tConsumedTimes: %d, FirstConsumeTime: %d, NextConsumeTime: %d\n"+
                            "\tBody: %s\n"+
                            "\tProps: %s\n",
                            v.MessageId, v.PublishTime, v.MessageTag, v.ConsumedTimes,
                            v.FirstConsumeTime, v.NextConsumeTime, v.MessageBody, v.Properties)
                    }

                    // If the broker does not receive an acknowledgment (ACK) before the NextConsumeTime
                    // elapses, the message is redelivered to the consumer.
                    // A unique timestamp is specified for the handle each time the message is consumed.
                    ackerr := mqConsumer.AckMessage(handles)
                    if ackerr != nil {
                        // If the handle times out, the broker fails to receive an ACK from the consumer.
                        fmt.Println(ackerr)
                        if errAckItems, ok := ackerr.(errors.ErrCode).Context()["Detail"].([]mq_http_sdk.ErrAckItem); ok {
                           for _, errAckItem := range errAckItems {
                             fmt.Printf("\tErrorHandle:%s, ErrorCode:%s, ErrorMsg:%s\n",
                               errAckItem.ErrorHandle, errAckItem.ErrorCode, errAckItem.ErrorMsg)
                           }
                        } else {
                           fmt.Println("ack err =", ackerr)
                        }
                        time.Sleep(time.Duration(3) * time.Second)
                    } else {
                        fmt.Printf("Ack ---->\n\t%s\n", handles)
                    }

                    endChan <- 1
                }
            case err := <-errChan:
                {
                    // No message is available for consumption in the topic.
                    if strings.Contains(err.(errors.ErrCode).Error(), "MessageNotExist") {
                        fmt.Println("\nNo new message, continue!")
                    } else {
                        fmt.Println(err)
                        time.Sleep(time.Duration(3) * time.Second)
                    }
                    endChan <- 1
                }
            case <-time.After(35 * time.Second):
                {
                    fmt.Println("Timeout of consumer message ??")
                    endChan <- 1
                }
            }
        }()

        // Long polling: the default network timeout is 35 seconds.
        // If no message is available, the request is suspended on the broker for the specified period.
        // If a message becomes available during the wait, the broker responds immediately.
        mqConsumer.ConsumeMessage(respChan, errChan,
            3, // Maximum number of messages per batch (max: 16).
            3, // Long polling wait time in seconds (max: 30).
        )
        <-endChan
    }
}

Étapes suivantes