Cette rubrique propose des exemples de code Java pour envoyer et consommer des messages planifiés ou différés via le SDK client HTTP d'ApsaraMQ for RocketMQ.
Messages planifiés et messages différés
Au niveau de l'API, les messages planifiés et les messages différés constituent une seule et même fonctionnalité. Tous deux utilisent la méthode setStartDeliverTime pour différer la livraison. La seule différence réside dans le calcul de l'horodatage :
| Type | Méthode de définition de l'heure de livraison | Exemple |
|---|---|---|
| Message différé | Heure actuelle + décalage temporel | System.currentTimeMillis() + 10 * 1000 (délai de 10 secondes) |
| Message planifié | Horodatage Unix absolu en millisecondes | 1654770600000 (2022-06-09 18:30:00) |
Le broker conserve chaque message jusqu'à l'heure de livraison spécifiée, puis le libère pour consommation. Ce paramètre s'applique par message, ce qui permet à chaque message d'avoir sa propre heure de livraison.
Pour plus d'informations, consultez la documentation Messages planifiés et messages différés.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Créé une instance, un topic et un groupe de consommateurs dans la console ApsaraMQ for RocketMQ
Créé une paire de clés AccessKey pour votre compte Alibaba Cloud
Chaque topic ne prend en charge qu'un seul type de message. N'utilisez pas un topic configuré pour des messages normaux pour envoyer des messages planifiés ou différés.
Envoyer des messages planifiés ou différés
Tous les messages de cet exemple utilisent la méthode setStartDeliverTime pour différer la livraison de 10 secondes. Pour envoyer un message planifié, remplacez le décalage relatif par un horodatage Unix absolu en millisecondes (par exemple, 1654770600000 pour le 09/06/2022 à 18:30:00).
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(
// HTTP endpoint. Find this on the Instance Details page, in the
// Endpoints tab under Basic Information.
"<your-http-endpoint>",
// Load credentials from environment variables to avoid hardcoding secrets.
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
);
final String topic = "<your-topic>";
// If the instance has a namespace, specify the instance ID.
// If it does not have a namespace, set this to null or "".
final String instanceId = "<your-instance-id>";
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
);
pubMsg.getProperties().put("a", String.valueOf(i)); // Custom property
pubMsg.setMessageKey("MessageKey");
// Defer delivery by 10 seconds.
// For a scheduled message, use an absolute timestamp instead,
// e.g., pubMsg.setStartDeliverTime(1654770600000L);
pubMsg.setStartDeliverTime(System.currentTimeMillis() + 10 * 1000);
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();
}
}
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Emplacement |
|---|---|---|
<your-http-endpoint> |
Endpoint HTTP | Page Instance Details > onglet Endpoints > HTTP Endpoint |
<your-topic> |
Nom du topic | Console ApsaraMQ for RocketMQ > Topics |
<your-instance-id> |
ID de l'instance, ou null si l'instance n'a pas de namespace |
Page Instance Details > section Basic Information |
Consommer des messages planifiés ou différés
Consommez les messages planifiés et différés de la même manière que les messages normaux. Le broker les conserve jusqu'à l'heure de livraison, puis les libère pour consommation.
L'exemple suivant utilise le long polling pour consommer les messages et renvoyer des accusés de réception (ACK) au broker.
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(
// HTTP endpoint. Find this on the Instance Details page, in the
// Endpoints tab under Basic Information.
"<your-http-endpoint>",
// The AccessKey ID for authentication.
"${ACCESS_KEY}",
// The AccessKey secret for authentication.
"${SECRET_KEY}"
);
final String topic = "<your-topic>";
final String groupId = "<your-group-id>";
// If the instance has a namespace, specify the instance ID.
// If it does not have a namespace, set this to null or "".
final String instanceId = "<your-instance-id>";
final MQConsumer consumer;
if (instanceId != null && instanceId != "") {
consumer = mqClient.getConsumer(instanceId, topic, groupId, null);
} else {
consumer = mqClient.getConsumer(topic, groupId);
}
// Consume messages in a loop. For production use, run multiple threads
// to consume messages concurrently.
do {
List<Message> messages = null;
try {
// Long polling: wait up to 3 seconds for new messages.
// - First parameter: max messages per batch (max: 16).
// - Second parameter: long polling timeout in seconds (max: 30).
messages = consumer.consumeMessage(3, 3);
} 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);
}
// ACK consumed messages. If the broker does not receive an ACK
// before the retry interval elapses, the message is redelivered.
{
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);
}
}
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Emplacement |
|---|---|---|
<your-http-endpoint> |
Endpoint HTTP | Page Instance Details > onglet Endpoints > HTTP Endpoint |
${ACCESS_KEY} |
AccessKey ID pour l'authentification | Consultez Create an AccessKey pair |
${SECRET_KEY} |
AccessKey secret pour l'authentification | Consultez Create an AccessKey pair |
<your-topic> |
Nom du topic | Console ApsaraMQ for RocketMQ > Topics |
<your-group-id> |
ID du groupe de consommateurs | Console ApsaraMQ for RocketMQ > Groups |
<your-instance-id> |
ID de l'instance, ou null si l'instance n'a pas de namespace |
Page Instance Details > section Basic Information |
Étapes suivantes
Messages planifiés et messages différés : Découvrez le cycle de vie, les règles de livraison et les limites des messages planifiés et différés.
Envoyer et recevoir des messages normaux : Envoyez des messages livrés immédiatement sans aucun délai.