Utilisez le SDK client HTTP Python pour envoyer et recevoir des messages planifiés et différés dans ApsaraMQ for RocketMQ. Cette rubrique fournit des exemples complets de code pour les producteurs et les consommateurs.
Fonctionnement des messages planifiés et différés
Les messages planifiés et différés retardent la livraison aux consommateurs jusqu’à une heure spécifiée :
Message différé : livré après un délai fixe. Par exemple, un message avec un délai de 10 secondes devient disponible pour les consommateurs 10 secondes après son envoi.
Message planifié : livré à un moment précis. Par exemple, un message planifié pour le 15/01/2024 à 15:00:00 est livré à cette heure précise.
Via HTTP, les deux types utilisent la même API. Définissez start_deliver_time sur un horodatage Unix futur en millisecondes :
Pour un message différé, calculez l’horodatage comme suit :
current_time + delay_duration.Pour un message planifié, définissez directement l’horodatage sur l’heure de livraison cible.
Pour plus d’informations, consultez la section Messages planifiés et messages différés.
Prérequis
Avant de commencer, assurez-vous d’avoir :
Le SDK client HTTP Python installé. Pour plus de détails, consultez la section Préparation de l'environnement
Une instance, une rubrique et un groupe de consommateurs ApsaraMQ for RocketMQ créés dans la console ApsaraMQ for RocketMQ. Pour plus de détails, consultez la section Création de ressources
Une paire AccessKey pour votre compte Alibaba Cloud. Pour plus de détails, consultez la section Création d'une paire AccessKey
Envoi de messages planifiés ou différés
L’exemple suivant envoie quatre messages différés, chacun étant livré 10 secondes après son envoi. Pour envoyer plutôt un message planifié, définissez start_deliver_time sur l’heure de livraison cible sous forme d’horodatage Unix en millisecondes.
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Exemple |
|---|---|---|
<http-endpoint> |
Endpoint HTTP de votre instance. Vous le trouverez sur la page Instance Details dans la console ApsaraMQ for RocketMQ. | http://xxxxx.mqrest.cn-hangzhou.aliyuncs.com |
<topic> |
Nom de la rubrique vers laquelle envoyer les messages. | scheduled-msg-topic |
<instance-id> |
ID de l’instance propriétaire de la rubrique. Requis si l’instance possède un namespace. Si l’instance n’a pas de namespace, transmettez une chaîne vide. Vérifiez le statut du namespace sur la page Instance Details. | MQ_INST_xxxxx |
import sys
import os
import time
from mq_http_sdk.mq_exception import MQExceptionBase
from mq_http_sdk.mq_producer import *
from mq_http_sdk.mq_client import *
# Initialize the producer client
mq_client = MQClient(
# HTTP endpoint. Find it in the HTTP Endpoint section
# of the Instance Details page in the ApsaraMQ for RocketMQ console.
"<http-endpoint>",
# AccessKey ID and AccessKey secret for authentication.
# Store credentials in environment variables instead of hardcoding them.
os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
)
topic_name = "<topic>"
instance_id = "<instance-id>"
producer = mq_client.get_producer(instance_id, topic_name)
# Send four messages
msg_count = 4
print("%sPublish Message To %s\nTopicName:%s\nMessageCount:%s\n" % (10 * "=", 10 * "=", topic_name, msg_count))
try:
for i in range(msg_count):
msg = TopicMessage(
"I am test message %s.hello" % i, # Message content
"tag1" # Message tag
)
# Set a custom message attribute
msg.put_property("a", "i")
# Set the message key for tracing
msg.set_message_key("MessageKey")
# Deliver this message after a 10-second delay.
# The value is a Unix timestamp in milliseconds.
# For a scheduled message, replace the calculation below
# with the target delivery time in milliseconds.
msg.set_start_deliver_time(int(round(time.time() * 1000)) + 10 * 1000)
re_msg = producer.publish_message(msg)
print("Publish Timer Message Succeed. MessageID:%s, BodyMD5:%s" % (re_msg.message_id, re_msg.message_body_md5))
except MQExceptionBase as e:
if e.type == "TopicNotExist":
print("Topic not exist, please create it.")
sys.exit(1)
print("Publish Message Fail. Exception:%s" % e)
**Paramètre clé : set_start_deliver_time**
Ce paramètre contrôle le moment où le broker livre le message. Transmettez un horodatage Unix en millisecondes :
| Type de livraison | Définition de l'horodatage | Exemple |
|---|---|---|
| Livraison différée | Heure actuelle plus la durée du délai | int(round(time.time() * 1000)) + 10 * 1000 |
| Livraison planifiée | Horodatage absolu de l’heure de livraison cible | 1700000000000 |
Réception de messages planifiés et différés
Le broker conserve les messages planifiés et différés jusqu’à l’arrivée de l’heure de livraison, puis les rend disponibles pour la consommation. Consommez-les de la même manière que les messages normaux.
L’exemple suivant utilise le long polling pour recevoir et acquitter les messages. Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Exemple |
|---|---|---|
<http-endpoint> |
Endpoint HTTP de votre instance. | http://xxxxx.mqrest.cn-hangzhou.aliyuncs.com |
<topic> |
Nom de la rubrique depuis laquelle consommer les messages. | scheduled-msg-topic |
<group-id> |
ID du groupe de consommateurs. | GID_scheduled_msg |
<instance-id> |
ID de l’instance. Transmettez une chaîne vide si l’instance ne possède pas de namespace. | MQ_INST_xxxxx |
import os
import time
from mq_http_sdk.mq_exception import MQExceptionBase
from mq_http_sdk.mq_consumer import *
from mq_http_sdk.mq_client import *
# Initialize the consumer client
mq_client = MQClient(
"<http-endpoint>",
os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID'],
os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
)
topic_name = "<topic>"
group_id = "<group-id>"
instance_id = "<instance-id>"
consumer = mq_client.get_consumer(instance_id, topic_name, group_id)
# Long polling timeout in seconds (max: 30)
wait_seconds = 3
# Maximum messages per request (max: 16)
batch = 3
print(("%sConsume And Ak Message From Topic%s\nTopicName:%s\nMQConsumer:%s\nWaitSeconds:%s\n" \
% (10 * "=", 10 * "=", topic_name, group_id, wait_seconds)))
while True:
try:
# Long polling: the request blocks on the broker until a message
# arrives or wait_seconds elapses, whichever comes first.
recv_msgs = consumer.consume_message(batch, wait_seconds)
for msg in recv_msgs:
print(("Receive, MessageId: %s\nMessageBodyMD5: %s \
\nMessageTag: %s\nConsumedTimes: %s \
\nPublishTime: %s\nBody: %s \
\nNextConsumeTime: %s \
\nReceiptHandle: %s \
\nProperties: %s\n" % \
(msg.message_id, msg.message_body_md5,
msg.message_tag, msg.consumed_times,
msg.publish_time, msg.message_body,
msg.next_consume_time, msg.receipt_handle, msg.properties)))
except MQExceptionBase as e:
if e.type == "MessageNotExist":
print(("No new message! RequestId: %s" % e.req_id))
continue
print(("Consume Message Fail! Exception:%s\n" % e))
time.sleep(2)
continue
# Acknowledge consumed messages.
# If the broker does not receive an ACK before msg.next_consume_time,
# it redelivers the message. Each delivery generates a new receipt handle.
try:
receipt_handle_list = [msg.receipt_handle for msg in recv_msgs]
consumer.ack_message(receipt_handle_list)
print(("Ak %s Message Succeed.\n\n" % len(receipt_handle_list)))
except MQExceptionBase as e:
print(("\nAk Message Fail! Exception:%s" % e))
if e.sub_errors:
for sub_error in e.sub_errors:
print(("\tErrorHandle:%s,ErrorCode:%s,ErrorMsg:%s" % \
(sub_error["ReceiptHandle"], sub_error["ErrorCode"], sub_error["ErrorMessage"])))
Fonctionnement du long polling
La méthode consume_message bloque la requête sur le broker pendant au maximum wait_seconds. Si un message devient disponible durant cette période, le broker répond immédiatement. Sinon, la requête expire après le délai imparti.
| Paramètre | Description | Plage |
|---|---|---|
wait_seconds |
Délai d’expiration du long polling en secondes. | Maximum : 30 |
batch |
Nombre maximal de messages renvoyés par requête. | Maximum : 16 |
Fonctionnement de l’accusé de réception des messages
Après avoir traité un message, acquittez-le en renvoyant le receipt handle au broker. Si aucun accusé de réception (ACK) n’est reçu avant next_consume_time, le broker relivre le message avec un nouveau receipt handle. Utilisez toujours le handle provenant de la dernière livraison.
Étapes suivantes
Messages planifiés et messages différés — concepts et détails des fonctionnalités
Préparation de l'environnement — configuration du SDK pour d’autres langages