ApsaraMQ for RocketMQ propose deux types de consommateurs : Push consumer et Simple consumer. Chaque type gère différemment la récupération des messages, la concurrence et les nouvelles tentatives. Choisissez le type qui correspond à votre modèle de traitement et à vos exigences de fiabilité.
Quel type de consommateur utiliser
| Scénario | Type recommandé | Raison |
|---|---|---|
| Temps de traitement prévisible, pas de threading personnalisé | Push consumer | Le SDK gère la récupération des messages, la concurrence et les nouvelles tentatives. Enregistrez un écouteur, traitez chaque message et renvoyez le résultat. |
| Temps de traitement variable, workflows personnalisés | Simple consumer | Votre application contrôle le moment de la récupération des messages, leur distribution entre les threads et l'accusé de réception une fois le traitement terminé. |
Le changement de type de consommateur n'affecte ni les ressources ApsaraMQ for RocketMQ existantes ni le traitement métier.
Étapes du traitement des messages
Les deux types de consommateurs suivent un cycle de vie en trois étapes :
Réception -- Récupération des messages depuis le serveur.
Traitement -- Exécution de la logique métier sur chaque message.
Validation -- Signalement du résultat (succès ou échec) au serveur.

Les deux types diffèrent par la manière dont chaque étape est gérée :
| Fonctionnalité | Push consumer | Simple consumer |
|---|---|---|
| Interface | Rappel d'écouteur -- implémentez la logique dans l'écouteur et renvoyez un résultat | L'application appelle les opérations API pour recevoir, traiter et acquitter les messages |
| Concurrence | Gérée par le SDK | Gérée par l'application |
| Flexibilité | Fortement encapsulé, moins flexible | Opérations atomiques, hautement personnalisable |
| Idéal pour | Consommation standard avec temps de traitement prévisible | Workflows personnalisés, distribution asynchrone ou consommation par lots |
| Classes SDK | PushConsumer, LitePushConsumer |
SimpleConsumer |
Push consumers
Un Push consumer encapsule la récupération des messages, la gestion des threads et la logique de nouvelle tentative. Enregistrez un écouteur de messages lors de l'initialisation et le SDK se charge du reste.
Fonctionnement
Le SDK utilise en interne un modèle de thread Reactor :
Un thread d'interrogation longue intégré extrait les messages du serveur de manière asynchrone.
Les messages sont placés dans une file d'attente de cache interne.
Le SDK distribue les messages aux threads de consommation, qui invoquent votre écouteur.

Résultats de l'écouteur
L'écouteur de messages doit renvoyer l'un des résultats suivants :
| Résultat | Constante Java SDK | Comportement |
|---|---|---|
| Succès | ConsumeResult.SUCCESS |
Le serveur met à jour la progression de la consommation. |
| Échec | ConsumeResult.FAILURE |
Le système effectue une nouvelle tentative selon la politique de nouvelle tentative de PushConsumer. |
| Exception levée | (traitée comme un échec) | Même comportement de nouvelle tentative qu'en cas d'échec explicite. |
Comportement en cas de délai dépassé
Si la logique de traitement bloque et empêche le message d'être traité dans le délai imparti, le SDK soumet de force un résultat d'échec et gère le message conformément à la politique de nouvelle tentative.
Un délai dépassé amène le SDK à soumettre un résultat d'échec, mais le thread de traitement actuel peut ne pas répondre à l'interruption et continuer à s'exécuter.
Contraintes de fiabilité
Le Push consumer détermine le succès ou l'échec strictement à partir de la valeur de retour de l'écouteur. Pour préserver cette garantie :
Traitez de manière synchrone. Terminez tout le traitement avant de renvoyer le résultat.
Ne redistribuez pas les messages. Ne confiez pas les messages à d'autres threads et ne renvoyez pas de résultat avant que ces threads n'aient terminé.
Si l'écouteur renvoie un succès avant la fin du traitement et que ce dernier échoue ensuite, le serveur considère le message comme consommé et n'effectue pas de nouvelle tentative.
Livraison des messages ordonnés
Lorsqu'un groupe de consommateurs utilise le mode de consommation ordonnée, le Push consumer invoque l'écouteur dans l'ordre strict des messages sans configuration supplémentaire. Pour plus d'informations, consultez Messages ordonnés.
La livraison ordonnée nécessite un traitement synchrone. Une distribution asynchrone personnalisée au sein de l'écouteur annule la garantie d'ordonnancement.
Quand utiliser les Push consumers
Les Push consumers sont idéaux lorsque :
Le temps de traitement est prévisible. Des durées imprévisibles déclenchent de fréquents dépassements de délai, entraînant des messages en double via des nouvelles tentatives.
Une consommation standard suffit. Le SDK contrôle le modèle de thread et livre les messages avec un débit maximal. Cela simplifie le développement mais ne prend pas en charge le traitement asynchrone ni le contrôle de débit personnalisé.
Classes SDK
ApsaraMQ for RocketMQ fournit deux classes SDK pour les Push consumers :
**
PushConsumer** -- Consomme des messages provenant de topics standards (non-Lite).**
LitePushConsumer** -- Consomme des messages provenant de topics de type Lite, avec un contrôle de la consommation à la granularité du topic Lite.
Exemple PushConsumer
// Consume normal messages with PushConsumer
ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "<your-topic>";
FilterExpression filterExpression = new FilterExpression("<your-filter-tag>", FilterExpressionType.TAG);
PushConsumer pushConsumer = provider.newPushConsumerBuilder()
// Set the consumer group
.setConsumerGroup("<your-consumer-group>")
// Set the access endpoint
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("<your-endpoint>").build())
// Bind the subscription
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
// Register the message listener
.setMessageListener(new MessageListener() {
@Override
public ConsumeResult consume(MessageView messageView) {
// Process the message and return the result
return ConsumeResult.SUCCESS;
}
})
.build();
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Exemple |
|---|---|---|
<your-topic> |
Nom du topic | order-events |
<your-filter-tag> |
Tag de message pour le filtrage | payment |
<your-consumer-group> |
Nom du groupe de consommateurs | order-service-group |
<your-endpoint> |
Endpoint d'accès au serveur | -- |
Exemple LitePushConsumer
// Consume normal messages with LitePushConsumer
ClientServiceProvider provider = ClientServiceProvider.loadService();
LitePushConsumer litePushConsumer = provider.newLitePushConsumerBuilder()
// Set the access endpoint
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("<your-endpoint>").build())
// Set the topic
.bindTopic("<your-topic>")
// Set the consumer group
.setConsumerGroup("<your-consumer-group>")
// Register the message listener
.setMessageListener(messageView -> {
// Process the message and return the result
return ConsumeResult.SUCCESS;
})
.build();
// Subscribe to Lite topics
litePushConsumer.subscribeLite("<your-lite-topic-1>");
litePushConsumer.subscribeLite("<your-lite-topic-2>");
Simple consumers
Un Simple consumer fournit des opérations API atomiques pour le traitement des messages. Votre application contrôle directement la récupération des messages, la gestion des threads et l'accusé de réception.
Fonctionnement
Appelez
ReceiveMessagepour extraire un lot de messages du serveur.Distribuez les messages à vos threads métier pour traitement.
Appelez
AckMessagepour chaque message traité avec succès.
En cas d'échec du traitement, n'envoyez pas d'accusé de réception. Le message redevient disponible après l'expiration de la durée d'invisibilité, ce qui déclenche une nouvelle tentative. Pour plus d'informations, consultez Politique de nouvelle tentative de SimpleConsumer.
Opérations API
| Opération | Objectif | Paramètres clés |
|---|---|---|
ReceiveMessage |
Extraire des messages du serveur | Taille du lot : nombre de messages par requête. Durée d'invisibilité du message : temps de traitement maximal avant la redistribution du message. |
AckMessage |
Accuser réception d'une consommation réussie | Aucun |
ChangeInvisibleDuration |
Prolonger le temps de traitement pour un message déjà reçu | Durée d'invisibilité du message : nouvelle valeur, généralement utilisée lorsque le traitement prend plus de temps que prévu initialement. |
Le serveur utilisant un stockage distribué, ReceiveMessage peut renvoyer un résultat vide même si des messages existent. Pour gérer cela, appelez à nouveau ReceiveMessage ou augmentez la concurrence des appels.
Gestion des échecs
Le tableau suivant décrit l'impact des différents scénarios d'échec sur la livraison des messages :
| Scénario d'échec | Comportement |
|---|---|
| Échec du traitement (aucun ACK envoyé) | Le message redevient visible après l'expiration de la durée d'invisibilité. Le serveur le redistribue pour une nouvelle tentative. |
| Dépassement de la durée d'invisibilité pendant le traitement | Identique à l'absence d'ACK -- le message redevient visible et est redistribué. Utilisez ChangeInvisibleDuration pour prolonger la fenêtre de traitement avant l'expiration de la durée. |
| Arrêt brutal du consommateur avant l'envoi de l'ACK | Le message est redistribué après l'expiration de la durée d'invisibilité. |
Livraison des messages ordonnés
Un Simple consumer traite les messages ordonnés selon l'ordre de stockage. Pour un groupe de messages devant rester séquentiels, le message suivant ne peut être récupéré tant que le précédent n'est pas traité.
Quand utiliser les Simple consumers
Les Simple consumers sont idéaux lorsque :
Le temps de traitement est imprévisible. Spécifiez une durée d'invisibilité initiale du message lors de l'appel à
ReceiveMessage, puis prolongez-la avecChangeInvisibleDurationsi nécessaire.Des workflows personnalisés sont requis. Le SDK n'impose aucun modèle de threading -- implémentez une distribution asynchrone, une consommation par lots ou tout autre modèle personnalisé.
Le contrôle du débit est important. Votre code décide quand et à quelle fréquence appeler
ReceiveMessage, offrant un contrôle direct sur le débit.
Classe SDK
ApsaraMQ for RocketMQ fournit une classe SDK pour les Simple consumers : SimpleConsumer. Cette classe ne peut pas consommer de messages provenant de topics Lite.
Exemple SimpleConsumer
// Consume normal messages with SimpleConsumer
ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "<your-topic>";
FilterExpression filterExpression = new FilterExpression("<your-filter-tag>", FilterExpressionType.TAG);
SimpleConsumer simpleConsumer = provider.newSimpleConsumerBuilder()
// Set the consumer group
.setConsumerGroup("<your-consumer-group>")
// Set the access endpoint
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("<your-endpoint>").build())
// Bind the subscription
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
.build();
try {
// Pull up to 10 messages, wait up to 30 seconds
List<MessageView> messageViewList = simpleConsumer.receive(10, Duration.ofSeconds(30));
messageViewList.forEach(messageView -> {
System.out.println(messageView);
// Acknowledge each message after successful processing
try {
simpleConsumer.ack(messageView);
} catch (ClientException e) {
e.printStackTrace();
}
});
} catch (ClientException e) {
// Handle failures such as throttling, then retry the receive call
e.printStackTrace();
}
Bonnes pratiques
Contrôle du temps de traitement pour les Push consumers
Maintenez le traitement des messages dans les limites du seuil de délai. Des dépassements de délai fréquents provoquent des nouvelles tentatives inutiles et des messages en double. Si votre application gère régulièrement des tâches de longue durée, basculez vers un Simple consumer et définissez une durée d'invisibilité des messages appropriée.