Les exemples de code PHP ci-dessous montrent comment envoyer et consommer des messages planifiés et différés avec le SDK client HTTP d'ApsaraMQ for RocketMQ.
Fonctionnement des messages planifiés et différés
Les messages planifiés et différés retardent leur livraison au lieu d'être remis immédiatement après publication.
Message différé : livré après un délai spécifié. Exemple : livraison 30 secondes après l'envoi.
Message planifié : livré à un moment précis. Exemple : livraison le 15/06/2024 à 09:00:00.
Ces deux types partagent la même API. Appelez setStartDeliverTime avec un horodatage Unix en millisecondes indiquant quand le broker doit remettre le message.
Pour un message différé, calculez l'horodatage comme suit :
heure actuelle + durée du délai.Pour un message planifié, convertissez directement l'heure cible de livraison en horodatage Unix en millisecondes.
Cas d'utilisation courants :
Expiration de commande : annulez une commande impayée après 30 minutes en publiant un message différé lors de sa création.
Nouvelle tentative avec backoff : retraitez une tâche ayant échoué après un délai plutôt que d'effectuer une nouvelle tentative immédiate.
Notifications programmées : déclenchez un rappel à une heure précise, par exemple 15 minutes avant un événement planifié.
Pour plus d'informations, consultez la rubrique Scheduled messages and delayed messages.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Créé les ressources dans la console ApsaraMQ for RocketMQ : une instance, un topic et un groupe de consommateurs
Créé une paire AccessKey et enregistré l'AccessKey ID et l'AccessKey secret en tant que variables d'environnement
Envoyer des messages planifiés ou différés
Le code suivant envoie quatre messages différés, chacun étant livré 10 secondes après publication. Pour envoyer un message planifié à la place, transmettez l'heure cible de livraison sous forme d'horodatage Unix en millisecondes à setStartDeliverTime.
Remplacez les espaces réservés avant d'exécuter le code :
| Espace réservé | Description | Exemple |
|---|---|---|
<HTTP_ENDPOINT> |
Endpoint HTTP issu de la page Instance Details dans la console ApsaraMQ for RocketMQ | http://xxxx.mq-http.cn-hangzhou.aliyuncs.com |
<TOPIC> |
Nom du topic créé dans la console | scheduled-msg-topic |
<INSTANCE_ID> |
ID de l'instance. Si l'instance n'a pas de namespace, définissez la valeur sur null ou "" |
MQ_INST_xxxx |
<?php
require "vendor/autoload.php";
use MQ\Model\TopicMessage;
use MQ\MQClient;
class ProducerTest
{
private $client;
private $producer;
public function __construct()
{
$this->client = new MQClient(
// HTTP endpoint from the Instance Details page in the ApsaraMQ for RocketMQ console
"<HTTP_ENDPOINT>",
// AccessKey ID for authentication
getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
// AccessKey secret for authentication
getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET')
);
$topic = "<TOPIC>";
// Instance ID. If the instance has no namespace, set to null or "".
$instanceId = "<INSTANCE_ID>";
$this->producer = $this->client->getProducer($instanceId, $topic);
}
public function run()
{
try
{
for ($i = 1; $i <= 4; $i++)
{
$publishMessage = new TopicMessage(
"hello mq! " // Message body
);
// Custom message property
$publishMessage->putProperty("a", $i);
// Message key for tracing
$publishMessage->setMessageKey("MessageKey");
// Deliver the message 10 seconds from now.
// For a scheduled message, set this to the target delivery time
// as a millisecond-level Unix timestamp.
$publishMessage->setStartDeliverTime(time() * 1000 + 10 * 1000);
$result = $this->producer->publishMessage($publishMessage);
print "Send mq message success. msgId is:" . $result->getMessageId()
. ", bodyMD5 is:" . $result->getMessageBodyMD5() . "\n";
}
} catch (\Exception $e) {
print_r($e->getMessage() . "\n");
}
}
}
$instance = new ProducerTest();
$instance->run();
?>
Points clés :
setStartDeliverTimeaccepte un horodatage Unix en millisecondes. En PHP,time()renvoie des secondes ; multipliez par 1000 pour convertir.Pour un délai de 10 secondes :
time() * 1000 + 10 * 1000.Pour une livraison à une heure précise, convertissez cette heure en horodatage Unix en millisecondes. Par exemple, pour livrer à
2024-06-15 09:00:00 UTC:strtotime('2024-06-15 09:00:00') * 1000.
Consommer des messages planifiés ou différés
Les messages planifiés et différés se consomment de la même manière que les messages normaux. Une fois l'heure de livraison atteinte, les messages deviennent disponibles via le long polling.
Remplacez les espaces réservés avant d'exécuter le code :
| Espace réservé | Description | Exemple |
|---|---|---|
<HTTP_ENDPOINT> |
Endpoint HTTP issu de la page Instance Details | http://xxxx.mq-http.cn-hangzhou.aliyuncs.com |
<TOPIC> |
Nom du topic créé dans la console | scheduled-msg-topic |
<GROUP_ID> |
ID du groupe de consommateurs créé dans la console | GID_scheduled_consumer |
<INSTANCE_ID> |
ID de l'instance. Si l'instance n'a pas de namespace, définissez la valeur sur null ou "" |
MQ_INST_xxxx |
<?php
use MQ\MQClient;
require "vendor/autoload.php";
class ConsumerTest
{
private $client;
private $consumer;
public function __construct()
{
$this->client = new MQClient(
// HTTP endpoint from the Instance Details page in the ApsaraMQ for RocketMQ console
"<HTTP_ENDPOINT>",
// AccessKey ID for authentication
getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
// AccessKey secret for authentication
getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET')
);
$topic = "<TOPIC>";
// Consumer group ID created in the ApsaraMQ for RocketMQ console
$groupId = "<GROUP_ID>";
// Instance ID. If the instance has no namespace, set to null or "".
$instanceId = "<INSTANCE_ID>";
$this->consumer = $this->client->getConsumer($instanceId, $topic, $groupId);
}
public function ackMessages($receiptHandles)
{
try {
$this->consumer->ackMessage($receiptHandles);
} catch (\Exception $e) {
if ($e instanceof MQ\Exception\AckMessageException) {
// ACK may fail if the receipt handle has expired
printf("Ack Error, RequestId:%s\n", $e->getRequestId());
foreach ($e->getAckMessageErrorItems() as $errorItem) {
printf("\tReceiptHandle:%s, ErrorCode:%s, ErrorMsg:%s\n",
$errorItem->getReceiptHandle(),
$errorItem->getErrorCode(),
$errorItem->getErrorMessage());
}
}
}
}
public function run()
{
// Consume messages in a loop. For production use, run multiple threads
// to consume messages concurrently.
while (True) {
try {
// Long polling: if no message is available, the request waits
// on the broker until a message arrives or the polling period ends.
$messages = $this->consumer->consumeMessage(
3, // Max messages per request (up to 16)
3 // Long polling timeout in seconds (up to 30)
);
} catch (\MQ\Exception\MessageResolveException $e) {
// Thrown when some messages contain invalid characters and cannot be parsed
$messages = $e->getPartialResult()->getMessages();
$failMessages = $e->getPartialResult()->getFailResolveMessages();
$receiptHandles = array();
foreach ($messages as $message) {
$receiptHandles[] = $message->getReceiptHandle();
printf("MsgID %s\n", $message->getMessageId());
}
foreach ($failMessages as $failMessage) {
$receiptHandles[] = $failMessage->getReceiptHandle();
printf("Fail To Resolve Message. MsgID %s\n", $failMessage->getMessageId());
}
$this->ackMessages($receiptHandles);
continue;
} catch (\Exception $e) {
if ($e instanceof MQ\Exception\MessageNotExistException) {
// No messages available. Long polling continues on the next iteration.
printf("No message, continue long polling! RequestId:%s\n", $e->getRequestId());
continue;
}
print_r($e->getMessage() . "\n");
sleep(3);
continue;
}
print "consume finish, messages:\n";
$receiptHandles = array();
foreach ($messages as $message) {
$receiptHandles[] = $message->getReceiptHandle();
printf("MessageID:%s TAG:%s BODY:%s \nPublishTime:%d, FirstConsumeTime:%d, \nConsumedTimes:%d, NextConsumeTime:%d, MessageKey:%s\n",
$message->getMessageId(), $message->getMessageTag(), $message->getMessageBody(),
$message->getPublishTime(), $message->getFirstConsumeTime(),
$message->getConsumedTimes(), $message->getNextConsumeTime(),
$message->getMessageKey());
print_r($message->getProperties());
}
// If the broker does not receive an ACK before the time specified
// in getNextConsumeTime(), the message is redelivered.
// Each delivery generates a new receipt handle.
print_r($receiptHandles);
$this->ackMessages($receiptHandles);
print "ack finish\n";
}
}
}
$instance = new ConsumerTest();
$instance->run();
?>
Informations importantes sur la consommation de messages
Long polling : lorsqu'aucun message n'est disponible, le broker conserve la requête du consommateur jusqu'à l'arrivée d'un message ou l'expiration du délai de polling. Cela réduit les requêtes réseau inutiles.
ACK (accusé de réception) : après avoir traité un message, le consommateur envoie un ACK avec le receipt handle. Si le broker ne reçoit aucun ACK avant
getNextConsumeTime(), il relivre le message.Expiration du receipt handle : chaque livraison attribue un receipt handle unique avec sa propre durée d'expiration. Un handle expiré entraîne l'échec de l'ACK et déclenche une nouvelle livraison.
Échecs d'analyse : si le corps du message contient des caractères non valides, le SDK lève une exception
MessageResolveException. AppelezgetPartialResult()pour récupérer les messages analysés avec succès et traiter les échecs séparément.