Cette rubrique propose des exemples Java pour envoyer et recevoir des messages normaux via le SDK client HTTP d'ApsaraMQ for RocketMQ. Les messages normaux constituent le type par défaut : ils ne comportent aucune sémantique de livraison particulière, contrairement aux messages planifiés, différés, ordonnés ou transactionnels, et sont transmis aux consommateurs aussi rapidement que possible.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Installé le SDK Java. Pour plus d'informations, consultez la section Préparer l'environnement
Créé les ressources requises dans la console ApsaraMQ for RocketMQ : une instance, un topic (configuré avec le type de message Normal) et un groupe de consommateurs. Pour plus d'informations, consultez la section Créer des ressources
Récupéré votre paire AccessKey (AccessKey ID et AccessKey secret). Pour plus d'informations, consultez la section Créer une paire AccessKey
Fonctionnement
Les deux exemples suivent le même flux de travail :
Créez un objet
MQClienten indiquant l'endpoint HTTP et les identifiants AccessKey.Associez un producteur ou un consommateur à une instance et à un topic spécifiques.
Envoyez ou interrogez les messages, puis fermez le client une fois l'opération terminée.
Chaque topic ne prend en charge qu'un seul type de message. Un topic créé pour les messages normaux ne peut ni envoyer ni recevoir de messages planifiés, différés, ordonnés ou transactionnels.
Envoyer un message
L'exemple suivant envoie quatre messages normaux vers un topic dans une boucle, en utilisant la publication synchrone. Si la méthode publishMessage se termine sans lever d'exception, le message a bien été envoyé.
Remplacez les espaces réservés par vos valeurs réelles :
| Espace réservé | Description | Où le trouver |
|---|---|---|
<your-http-endpoint> |
Endpoint HTTP de l'instance | Page Instance Details > Section HTTP Endpoint |
<your-topic> |
Nom du topic | Console ApsaraMQ for RocketMQ |
<your-instance-id> |
ID de l'instance. Définissez cette valeur sur null ou "" si l'instance n'utilise pas de namespace |
Page Instance Details |
import com.aliyun.mq.http.MQClient;
import com.aliyun.mq.http.MQProducer;
import com.aliyun.mq.http.model.TopicMessage;
import java.util.Date;
public class Producer {
public static void main(String[] args) {
MQClient mqClient = new MQClient(
"<your-http-endpoint>",
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
);
final String topic = "<your-topic>";
final String instanceId = "<your-instance-id>";
// Get a producer for this topic and instance.
MQProducer producer;
if (instanceId != null && instanceId != "") {
producer = mqClient.getProducer(instanceId, topic);
} else {
producer = mqClient.getProducer(topic);
}
try {
for (int i = 0; i < 4; i++) {
TopicMessage pubMsg = new TopicMessage(
"hello mq!".getBytes(), // Message body
"A" // Message tag for filtering
);
// Custom property
pubMsg.getProperties().put("a", String.valueOf(i));
// Business key for message tracing (use an order ID, user ID, or similar identifier)
pubMsg.setMessageKey("MessageKey");
// Publish synchronously. No exception means success.
TopicMessage pubResultMsg = producer.publishMessage(pubMsg);
System.out.println(new Date() + " Send mq message success. Topic is:" + topic
+ ", msgId is: " + pubResultMsg.getMessageId()
+ ", bodyMD5 is: " + pubResultMsg.getMessageBodyMD5());
}
} catch (Throwable e) {
System.out.println(new Date() + " Send mq message failed. Topic is:" + topic);
e.printStackTrace();
}
mqClient.close();
}
}
Points clés :
Identifiants issus des variables d'environnement : le code lit les valeurs
ALIBABA_CLOUD_ACCESS_KEY_IDetALIBABA_CLOUD_ACCESS_KEY_SECRETdepuis les variables d'environnement. Définissez ces variables avant d'exécuter le code.Tag de message : les tags permettent aux abonnés de filtrer les messages au sein d'un topic. Dans cet exemple, le tag est
"A".Clé de message : définissez la clé de message à l'aide d'un identifiant métier, tel qu'un ID de commande ou un ID utilisateur. Cela facilite la recherche et le suivi des messages dans la console ApsaraMQ for RocketMQ.
Recevoir et acquitter des messages
L'exemple suivant consomme les messages d'un topic dans une boucle utilisant le long polling. Après le traitement de chaque lot, il envoie un accusé de réception (ACK) au broker. Si le broker ne reçoit pas d'ACK avant l'expiration de l'intervalle de nouvelle tentative de livraison, il redistribue le message.
Remplacez les espaces réservés par vos valeurs réelles :
| Espace réservé | Description | Où le trouver |
|---|---|---|
<your-http-endpoint> |
Endpoint HTTP de l'instance | Page Instance Details > Section HTTP Endpoint |
<your-topic> |
Nom du topic | Console ApsaraMQ for RocketMQ |
<your-instance-id> |
ID de l'instance. Définissez cette valeur sur null ou "" si l'instance n'utilise pas de namespace |
Page Instance Details |
<your-group-id> |
ID du groupe de consommateurs | Console ApsaraMQ for RocketMQ |
import com.aliyun.mq.http.MQClient;
import com.aliyun.mq.http.MQConsumer;
import com.aliyun.mq.http.common.AckMessageException;
import com.aliyun.mq.http.model.Message;
import java.util.ArrayList;
import java.util.List;
public class Consumer {
public static void main(String[] args) {
MQClient mqClient = new MQClient(
"<your-http-endpoint>",
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
);
final String topic = "<your-topic>";
final String groupId = "<your-group-id>";
final String instanceId = "<your-instance-id>";
// Get a consumer for this topic, group, and instance.
final MQConsumer consumer;
if (instanceId != null && instanceId != "") {
consumer = mqClient.getConsumer(instanceId, topic, groupId, null);
} else {
consumer = mqClient.getConsumer(topic, groupId);
}
// Poll for messages in a loop.
// In production, use multiple threads for concurrent consumption.
do {
List<Message> messages = null;
try {
// Long polling: if no message is available, the request hangs
// on the broker for up to the specified number of seconds.
messages = consumer.consumeMessage(
3, // Max messages per batch (up to 16)
3 // Long-polling timeout in seconds (up to 30)
);
} catch (Throwable e) {
e.printStackTrace();
try {
Thread.sleep(2000);
} catch (InterruptedException e1) {
e1.printStackTrace();
}
}
if (messages == null || messages.isEmpty()) {
System.out.println(Thread.currentThread().getName() + ": no new message, continue!");
continue;
}
// --- Process messages ---
for (Message message : messages) {
System.out.println("Receive message: " + message);
}
// --- Acknowledge messages ---
// Each message has a unique receipt handle that expires after
// the delivery retry interval. ACK before it expires to prevent
// redelivery.
{
List<String> handles = new ArrayList<String>();
for (Message message : messages) {
handles.add(message.getReceiptHandle());
}
try {
consumer.ackMessage(handles);
} catch (Throwable e) {
if (e instanceof AckMessageException) {
AckMessageException errors = (AckMessageException) e;
System.out.println("Ack message fail, requestId is:" + errors.getRequestId()
+ ", fail handles:");
if (errors.getErrorMessages() != null) {
for (String errorHandle : errors.getErrorMessages().keySet()) {
System.out.println("Handle:" + errorHandle
+ ", ErrorCode:" + errors.getErrorMessages().get(errorHandle).getErrorCode()
+ ", ErrorMsg:" + errors.getErrorMessages().get(errorHandle).getErrorMessage());
}
}
continue;
}
e.printStackTrace();
}
}
} while (true);
}
}
Points clés :
Long polling : le consommateur maintient la connexion ouverte sur le broker pendant la durée du délai d'attente spécifié (3 secondes dans cet exemple). Si un message arrive durant cette fenêtre, le broker répond immédiatement au lieu d'attendre le prochain cycle d'interrogation. Le délai maximal est de 30 secondes.
Taille du lot : chaque appel à
consumeMessagerenvoie jusqu'au nombre spécifié de messages (3 dans cet exemple, maximum 16).Receipt handle : à chaque livraison, un message reçoit un nouveau receipt handle associé à un horodatage unique. Utilisez ce handle pour acquitter le message. Si le handle expire avant que l'ACK n'atteigne le broker, le message est redistribué.
Concurrence : cet exemple utilise un seul thread. Pour obtenir un débit plus élevé, consommez les messages via plusieurs threads en environnement de production.
Remarques d'utilisation
Un seul type par topic : un topic créé pour les messages normaux ne peut ni envoyer ni recevoir d'autres types de messages (planifiés, différés, ordonnés ou transactionnels). Créez des topics distincts pour chaque type de message.
Définir des clés de message pertinentes : utilisez des identifiants métiers (tels que des IDs de commande ou des IDs utilisateur) comme clé de message. Cela accélère la recherche et le suivi des messages dans la console.
Gérer correctement les échecs d'acquittement : si
ackMessagelève une exceptionAckMessageException, enregistrez les handles ayant échoué ainsi que les codes d'erreur. Le broker redistribue automatiquement les messages non acquittés ; évitez donc de traiter les échecs d'ACK comme des erreurs fatales.Libérer les ressources : appelez
mqClient.close()lorsque le producteur ou le consommateur n'est plus nécessaire afin de libérer les connexions HTTP.
Étapes suivantes
Préparer l'environnement : installez le SDK Java et configurez les dépendances.
Créer des ressources : configurez les instances, les topics et les groupes de consommateurs dans la console.