Les messages normaux constituent le type de message par défaut dans ApsaraMQ for RocketMQ et ne reposent sur aucune sémantique de distribution particulière. Contrairement aux messages planifiés, différés, ordonnés ou transactionnels, les messages normaux n'embarquent aucun comportement de distribution supplémentaire.
Cette rubrique propose des exemples de code illustrant l'envoi et la réception de messages normaux à l'aide du SDK client HTTP pour Go.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Installé le SDK client HTTP pour Go. Pour plus d'informations, consultez la page Préparer l'environnement
Créé une instance, un topic et un groupe de consommateurs ApsaraMQ for RocketMQ dans la console. Pour plus d'informations, consultez la page Créer des ressources
Généré une paire AccessKey pour votre compte Alibaba Cloud. Pour plus d'informations, consultez la page Créer une paire AccessKey
Paramètres de configuration
Les deux exemples nécessitent les paramètres suivants. Remplacez les espaces réservés par vos valeurs réelles avant d'exécuter le code.
| Espace réservé | Description | Où le trouver |
|---|---|---|
<your-http-endpoint> |
Endpoint HTTP de votre instance | Page Instance Details > section HTTP Endpoint dans la console ApsaraMQ for RocketMQ |
<your-topic> |
Topic vers lequel envoyer les messages | Doit être créé dans la console. Chaque topic ne prend en charge qu'un seul type de message |
<your-instance-id> |
ID de votre instance | Page Instance Details. Si l'instance ne dispose pas de namespace, définissez cette valeur sur null ou sur une chaîne vide |
<your-group-id> |
ID de votre groupe de consommateurs (consommation uniquement) | Doit être créé dans la console |
Les identifiants AccessKey sont lus depuis les variables d'environnement :
ALIBABA_CLOUD_ACCESS_KEY_IDALIBABA_CLOUD_ACCESS_KEY_SECRET
Envoyer des messages normaux
L'exemple suivant initialise un producteur et envoie quatre messages vers un topic. Chaque message inclut un corps, un tag facultatif, une clé de message ainsi que des propriétés personnalisées.
package main
import (
"fmt"
"os"
"strconv"
"time"
"github.com/aliyunmq/mq-http-go-sdk"
)
func main() {
endpoint := "<your-http-endpoint>"
accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
topic := "<your-topic>"
instanceId := "<your-instance-id>"
client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")
producer := client.GetProducer(instanceId, topic)
for i := 0; i < 4; i++ {
msg := mq_http_sdk.PublishMessageRequest{
MessageBody: "hello mq!",
MessageTag: "", // Optional: filter tag
Properties: map[string]string{}, // Optional: custom properties
}
msg.MessageKey = "MessageKey"
msg.Properties["a"] = strconv.Itoa(i)
ret, err := producer.PublishMessage(msg)
if err != nil {
fmt.Println(err)
return
}
fmt.Printf("Publish ---->\n\tMessageId:%s, BodyMD5:%s, \n",
ret.MessageId, ret.MessageBodyMD5)
time.Sleep(100 * time.Millisecond)
}
}
Points clés :
PublishMessageest un appel synchrone qui renvoieMessageIdetMessageBodyMD5en cas de succès.Définissez
MessageTagpour acheminer les messages vers des consommateurs spécifiques abonnés avec un filtre de tag correspondant.Utilisez
MessageKeypour définir un identifiant personnalisé pour le message.
Consommer des messages normaux
L'exemple suivant initialise un consommateur et interroge continuellement les messages via un sondage long (long polling). Après le traitement de chaque lot, il accuse réception des messages afin d'éviter toute redistribution.
package main
import (
"fmt"
"os"
"strings"
"time"
"github.com/aliyunmq/mq-http-go-sdk"
"github.com/gogap/errors"
)
func main() {
endpoint := "<your-http-endpoint>"
accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
topic := "<your-topic>"
instanceId := "<your-instance-id>"
groupId := "<your-group-id>"
client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")
consumer := 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:
// Process the batch
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)
}
// Acknowledge processed messages
ackerr := consumer.AckMessage(handles)
if ackerr != nil {
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(3 * time.Second)
} else {
fmt.Printf("Ack ---->\n\t%s\n", handles)
}
endChan <- 1
case err := <-errChan:
if strings.Contains(err.(errors.ErrCode).Error(), "MessageNotExist") {
fmt.Println("\nNo new message, continue!")
} else {
fmt.Println(err)
time.Sleep(3 * time.Second)
}
endChan <- 1
case <-time.After(35 * time.Second):
fmt.Println("Timeout of consumer message ??")
endChan <- 1
}
}()
// Long polling: wait up to 3 seconds for new messages, fetch up to 3 per batch
consumer.ConsumeMessage(respChan, errChan,
3, // Max messages per batch (up to 16)
3, // Long polling wait time in seconds (up to 30)
)
<-endChan
}
}
Fonctionnement du sondage long (long polling)
En mode sondage long, si aucun message n'est disponible, la requête reste en attente sur le broker pendant la durée que vous spécifiez (jusqu'à 30 secondes). Dès qu'un message arrive, le broker répond immédiatement. Cette approche réduit le nombre de réponses vides par rapport au sondage court. Le délai d'expiration réseau par défaut est de 35 secondes.
Accusé de réception et redistribution
Après avoir traité un message, appelez AckMessage avec le handle de réception pour confirmer sa consommation. Si le broker ne reçoit pas d'accusé de réception avant l'échéance définie par NextConsumeTime, il redistribue le message. Chaque redistribution attribue un nouveau handle de réception.
Remarque : Si le handle de réception expire avant que vous n'accusiez réception du message, l'accusé échoue. La réponse d'erreur inclut ErrorHandle, ErrorCode et ErrorMsg pour chaque handle ayant échoué.