ApsaraMQ for RocketMQ prend en charge deux types de livraison de messages basée sur le temps : les messages planifiés, livrés à un instant précis, et les messages différés, livrés après un délai fixe. Ces deux types utilisent la propriété __STARTDELIVERTIME pour contrôler le moment où le broker livre le message aux consommateurs.
L'exemple de code Java ci-dessous montre comment envoyer des messages planifiés et différés, ainsi que s'y abonner, avec le SDK client TCP (Community Edition).
Fonctionnement
Lorsqu'un producteur envoie un message avec la propriété __STARTDELIVERTIME, le broker ApsaraMQ for RocketMQ conserve le message jusqu'à l'heure de livraison spécifiée. À ce moment-là, le broker livre le message aux consommateurs abonnés comme n'importe quel message normal.
Message différé : Définissez
__STARTDELIVERTIMEsurSystem.currentTimeMillis() + delayMillis. Le broker livre le message une fois le délai écoulé.Message planifié : Définissez
__STARTDELIVERTIMEsur un horodatage Unix futur en millisecondes. Le broker livre le message exactement à cet instant.
Si l'heure spécifiée est dans le passé, le broker livre immédiatement le message.
Planifié ou différé : quand utiliser chacun
| Type | Cas d'utilisation | Exemple |
|---|---|---|
| Différé | Déclencher une action après une attente fixe | Annuler une commande impayée 30 minutes après sa création |
| Planifié | Déclencher une action à une heure précise | Envoyer une notification à 09:00 chaque lundi |
Différences par rapport à Apache RocketMQ
Apache RocketMQ prend en charge les messages différés, mais pas les messages planifiés ; aucune interface dédiée n'existe pour ces derniers.
ApsaraMQ for RocketMQ offre des fonctionnalités supplémentaires :
Prise en charge des messages différés et planifiés
Précision à la seconde près pour l'heure de livraison
Concurrence plus élevée pour le traitement des messages basés sur le temps
Les méthodes de configuration et les résultats diffèrent entre Apache RocketMQ et ApsaraMQ for RocketMQ. Utilisez l'exemple de code de cette page pour ApsaraMQ for RocketMQ sur le cloud.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Le SDK Community Edition pour Java version 4.5.2 ou ultérieure installé
Une paire AccessKey créée pour votre compte Alibaba Cloud
Un topic et un ID de groupe créés dans la console ApsaraMQ for RocketMQ
Envoi de messages planifiés et différés
Les messages planifiés et différés utilisent tous deux la même propriété utilisateur __STARTDELIVERTIME. La seule différence réside dans le calcul de la valeur de l'horodatage.
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Exemple |
|---|---|---|
<your-group-id> |
ID de groupe créé dans la console ApsaraMQ for RocketMQ | GID_example |
<your-access-point> |
Endpoint de l'instance depuis la console | http://MQ_INST_XXXX.aliyuncs.com:80 |
<your-topic> |
Topic créé dans la console ApsaraMQ for RocketMQ | Topic_example |
<your-message-tag> |
Tag de message pour le filtrage | TagA |
import java.text.SimpleDateFormat;
import java.util.Date;
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;
public class RocketMQProducer {
/**
* Reads credentials from environment variables:
* ALIBABA_CLOUD_ACCESS_KEY_ID
* ALIBABA_CLOUD_ACCESS_KEY_SECRET
*/
private static RPCHook getAclRPCHook() {
return new AclClientRPCHook(new SessionCredentials(
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")));
}
public static void main(String[] args) throws MQClientException {
// Create a producer with message trace enabled.
// To disable message trace, use:
// new DefaultMQProducer("<your-group-id>", getAclRPCHook());
DefaultMQProducer producer = new DefaultMQProducer(
"<your-group-id>", getAclRPCHook(), true, null);
// Required for message trace on the cloud.
producer.setAccessChannel(AccessChannel.CLOUD);
// Set the instance endpoint from the ApsaraMQ for RocketMQ console.
producer.setNamesrvAddr("<your-access-point>");
producer.start();
for (int i = 0; i < 128; i++) {
try {
Message msg = new Message(
"<your-topic>",
"<your-message-tag>",
"Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));
// --- Option A: Delayed message ---
// Deliver 3 seconds from now.
long delayTime = System.currentTimeMillis() + 3000;
msg.putUserProperty("__STARTDELIVERTIME", String.valueOf(delayTime));
// --- Option B: Scheduled message ---
// Deliver at a specific date and time.
// Uncomment the following lines and comment out Option A to use.
//
// long timeStamp = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
// .parse("2025-08-10 18:45:00").getTime();
// msg.putUserProperty("__STARTDELIVERTIME", String.valueOf(timeStamp));
SendResult sendResult = producer.send(msg);
System.out.printf("%s%n", sendResult);
} catch (Exception e) {
System.out.println(new Date() + " Send mq message failed.");
e.printStackTrace();
}
}
// Shut down the producer before exiting (optional).
producer.shutdown();
}
}
Résumé :
La propriété
__STARTDELIVERTIMEaccepte un horodatage Unix en millisecondes.Pour les messages différés, ajoutez le délai souhaité (en millisecondes) à l'heure actuelle.
Pour les messages planifiés, convertissez la chaîne de date-heure cible au format
yyyy-MM-dd HH:mm:ssen horodatage Unix.Si l'heure spécifiée est antérieure à l'heure actuelle, le message est livré immédiatement.
Abonnement aux messages planifiés et différés
L'abonnement aux messages planifiés et différés est identique à celui des messages normaux. Le broker conserve le message jusqu'à l'heure de livraison spécifiée, donc aucune configuration spéciale côté consommateur n'est nécessaire.
import java.util.List;
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.client.consumer.rebalance.AllocateMessageQueueAveragely;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.RPCHook;
public class RocketMQPushConsumer {
/**
* Reads credentials from environment variables:
* ALIBABA_CLOUD_ACCESS_KEY_ID
* ALIBABA_CLOUD_ACCESS_KEY_SECRET
*/
private static RPCHook getAclRPCHook() {
return new AclClientRPCHook(new SessionCredentials(
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")));
}
public static void main(String[] args) throws MQClientException {
// Create a consumer with message trace enabled.
// To disable message trace, use:
// new DefaultMQPushConsumer("<your-group-id>", getAclRPCHook(),
// new AllocateMessageQueueAveragely());
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(
"<your-group-id>", getAclRPCHook(),
new AllocateMessageQueueAveragely(), true, null);
// Set the instance endpoint from the ApsaraMQ for RocketMQ console.
// The value is in the format of http://xxxx.mq-internet.aliyuncs.com:80.
consumer.setNamesrvAddr("<your-access-point>");
// Required for message trace on the cloud.
consumer.setAccessChannel(AccessChannel.CLOUD);
// Subscribe to all tags (*) on the topic.
consumer.subscribe("<your-topic>", "*");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(
List<MessageExt> msgs,
ConsumeConcurrentlyContext context) {
System.out.printf("Receive New Messages: %s %n", msgs);
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
}
}
Voir aussi
Messages planifiés et messages différés : concepts, précision de livraison et limites d'utilisation
Découvrez d'autres types de messages, tels que les messages transactionnels et les messages ordonnés