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 :
Installé le SDK HTTP Go. Pour plus de détails, consultez la page Préparer l'environnement
Créé une instance ApsaraMQ for RocketMQ, un topic et un groupe de consommateurs dans la console ApsaraMQ for RocketMQ
Obtenu une paire AccessKey pour votre compte Alibaba Cloud. Pour plus de détails, consultez la page Créer une paire AccessKey
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
Découvrez les concepts et les contraintes liés aux messages planifiés et aux messages différés
Créez des ressources telles que des instances, des topics et des groupes de consommateurs