Cette rubrique décrit la définition, les relations entre modèles, les attributs internes, les contraintes comportementales, la compatibilité des versions et les recommandations d'utilisation pour LiteTopic dans ApsaraMQ for RocketMQ.
Prérequis
Les sujets Lite sont pris en charge uniquement par les instances non Serverless (abonnement et paiement à l'utilisation) ainsi que par les instances Serverless dédiées.
-
Pour acheter une instance compatible avec le modèle de sujet Lite :
-
Lors de l'achat d'une nouvelle instance, ajoutez un tag de capacité produit sur la page d'achat avec la clé de tag version_capability et la valeur de tag lite-topic.
Pour les instances existantes, soumettez un ticket afin de mettre à niveau vers une version prenant en charge le modèle de sujet Lite. Indiquez l'ID de l'instance et la région lors de la soumission du ticket.
-
Soumettez un ticket pour obtenir une consultation gratuite sur les solutions LiteTopic adaptées à vos scénarios spécifiques.
Définition
Un sujet Lite est un conteneur secondaire pour la transmission et le stockage des messages dans ApsaraMQ for RocketMQ. Il sert à identifier les messages appartenant à différentes sous-classes (telles que différentes sessions, tâches ou autres granularités) au sein d'une même logique métier.
Les objectifs principaux des sujets Lite sont les suivants :
Permettre une consommation exclusive et définir une isolation des données de second niveau.
Nous vous recommandons de répartir les données issues de différentes sous-catégories dans des sujets Lite distincts afin d'obtenir une isolation plus fine du stockage et des abonnements.
Définir l'identité des données et les autorisations.
S'appuyant sur la gestion des identités et des autorisations basée sur les sujets, les sujets Lite permettent d'affiner davantage les identités utilisateur et les droits d'accès.
Relations entre modèles
Dans le modèle de domaine d'ApsaraMQ for RocketMQ, le flux et la position des sujets Lite sont les suivants :

Un sujet constitue le conteneur de premier niveau pour la transmission et le stockage des messages dans ApsaraMQ for RocketMQ. Lorsque le type de sujet est Lite, vous pouvez créer un sujet Lite sous celui-ci. La combinaison du sujet et du sujet Lite identifie de manière unique le conteneur de stockage des messages.
Lorsque le type de sujet est Lite, chaque conteneur de stockage utilise par défaut une seule file d'attente.
Propriétés internes
Nom du sujet Lite
Définition : Nom identifiant un sujet Lite. Les noms de sujets Lite sont globalement uniques au sein de leur sujet parent.
Valeur : Lorsque le type de sujet est Lite et que vous appelez setLiteTopic sur un message, le système crée automatiquement le sujet Lite s'il n'existe pas déjà.
Contrainte : Consultez les Limites des paramètres.
Durée de vie (TTL)
Définition : Durée d'expiration d'un sujet Lite. Si aucun nouveau message n'est écrit dans le sujet Lite pendant une période supérieure à sa TTL, le système le supprime automatiquement. La suppression libère le quota alloué au sujet Lite (le compteur total est décrémenté de un).
Valeur : Lors de la création d'un sujet de type Lite, vous pouvez définir le paramètre d'expiration.
Contrainte : Consultez les Limites des paramètres.
Compatibilité des versions
Version côté serveur : 5.0-rmq-20251024-1 ou ultérieure
Version côté client : RocketMQ gRPC 5.1.0 ou ultérieure
Différences entre les sujets de type Lite et de type Standard
|
Scénario |
Élément de comparaison |
Sujets légers |
Sujet de type standard |
|
Stockage des messages |
Sujet de premier niveau |
Identique. Les deux nécessitent la création préalable de la ressource de sujet. |
|
|
Sujet de second niveau |
Vous pouvez créer des millions de ressources LiteTopic de second niveau sous un sujet, chacune disposant de nouvelles fonctionnalités. |
Aucune ressource de sujet de second niveau. |
|
|
Gestion automatisée du cycle de vie |
La gestion du cycle de vie pour LiteTopic est automatisée :
|
Aucune |
|
|
Ordonnancement |
Chaque LiteTopic possède exactement une file d'attente. Les messages d'une même file d'attente sont stockés dans l'ordre.
|
Plusieurs files d'attente sont créées. Seuls les sujets ordonnés partitionnés garantissent l'ordonnancement. |
|
|
TPS maximal simultané pour l'envoi et la réception |
Étant donné que chaque LiteTopic ne dispose que d'une seule file d'attente, son TPS est limité. Toutefois, vous pouvez créer des millions de LiteTopics sous un seul sujet, de sorte que le TPS total évolue proportionnellement au nombre de LiteTopics. |
Le TPS du sujet évolue horizontalement en fonction du nombre de files d'attente et du nombre de nœuds du cluster. |
|
|
Consommation des messages |
Cohérence de l'abonnement |
Elles n'ont pas besoin d'être identiques. Au sein d'un même groupe, chaque consommateur peut s'abonner à un ensemble différent de LiteTopics. Les restrictions au niveau du groupe sont assouplies. |
Requise. Tous les consommateurs d'un même groupe doivent maintenir des abonnements identiques pour partager les messages du sujet cible. |
|
Ordonnancement |
Consommation ordonnée : les messages d'un LiteTopic sont traités par un seul thread de consommateur. |
Prend en charge la consommation simultanée ou ordonnée. |
|
|
Abonnement dynamique |
Chaque consommateur peut ajouter ou supprimer dynamiquement des abonnements à des LiteTopics spécifiques. |
Aucune |
|
|
Nombre maximal de LiteTopics auxquels un seul consommateur peut s'abonner |
Chaque consommateur peut s'abonner à des milliers de LiteTopics. |
Aucune |
|
|
Observabilité |
Métriques |
Comprend les métriques d'accumulation des messages. Aucune métrique n'est disponible pour le temps de latence du traitement des messages. |
Comprend les métriques d'accumulation des messages. Métrique du temps de traitement des messages |
|
Trace des messages |
Identique |
||
Cas d'utilisation courants des sujets Lite
Cas d'utilisation 1 : Communication asynchrone pour les systèmes Multi-Agent afin de résoudre le blocage des appels de longue durée
À mesure que les scénarios d'IA se complexifient, les systèmes mono-agent rencontrent des limites : manque de spécialisation, difficulté à intégrer plusieurs domaines et incapacité à permettre une prise de décision collaborative dynamique. Les applications et workflows mono-agent évoluent vers des architectures Multi-Agent. Toutefois, comme les tâches d'IA prennent souvent beaucoup de temps, les appels synchrones bloquent le thread de l'appelant, ce qui limite l'évolutivité pour une collaboration à grande échelle.

Comme illustré ci-dessus, le workflow Multi-Agent fonctionne de la manière suivante : l'agent Superviseur divise une requête en deux sous-tâches confiées à deux agents enfants. Chaque agent enfant résout sa partie et renvoie les résultats à l'agent Superviseur, qui les agrège et envoie la réponse finale au client web. L'utilisation de RocketMQ pour la communication asynchrone permet :
-
Flux de gestion des requêtes :
Créez un sujet (Request) pour chaque agent enfant servant de file d'attente tampon pour les tâches. Utilisez un sujet prioritaire pour traiter d'abord les tâches à haute priorité.
L'agent Superviseur envoie les détails des tâches divisées vers le sujet de requête correspondant.
-
Flux de gestion des réponses :
L'agent Superviseur crée un sujet de type Lite (Response) et s'y abonne.
Chaque agent enfant envoie le résultat de sa tâche vers un LiteTopic sous le sujet Response. Nommez chaque LiteTopic à l'aide de l'ID de tâche afin d'attribuer à chaque tâche son propre LiteTopic dédié.
L'agent Superviseur reçoit les résultats en temps réel via l'abonnement et les transmet au client web en utilisant HTTP SSE.
Cas d'utilisation 2 : Gestion distribuée de l'état des sessions pour résoudre les problèmes de continuité de session dans les applications d'IA
Les interactions des applications d'IA présentent des caractéristiques uniques : elles sont de longue durée, comportent plusieurs tours et dépendent fortement de ressources de calcul coûteuses par session. Lorsque les applications s'appuient sur des connexions persistantes telles que SSE, toute déconnexion (due à des redémarrages de passerelle, des délais d'expiration ou une instabilité du réseau) entraîne la perte du contexte de session actuel et gaspille les ressources de calcul d'IA déjà investies.

Comme dans le flux de réponse du cas d'utilisation 1, utilisez un sujet de type Lite pour les notifications de résultats en temps réel. Nommez chaque LiteTopic à l'aide du SessionID (par exemple, chatbot/{sessionID}). Tous les résultats de session sont livrés sous forme de messages ordonnés dans ce sujet. Pour maintenir la continuité de la session après une reconnexion :
Le client web établit une connexion persistante avec le nœud 1 du serveur d'application et démarre la session Session2.
Le nœud 1 du serveur d'application s'abonne au LiteTopic [chat/SessionID2].
Le planificateur de tâches du grand modèle de langage (LLM) envoie les résultats vers le LiteTopic [chat/SessionID2] en fonction du SessionID présent dans la requête.
En raison d'un problème réseau, le WebSocket se reconnecte au nœud 2 du serveur d'application.
Le nœud 1 du serveur d'application se désabonne du LiteTopic [chat/SessionID2]. Le nœud 2 du serveur d'application s'y abonne.
Le LiteTopic [chat/SessionID2] reprend la livraison à partir du dernier offset consommé, garantissant ainsi la continuité de l'état et des données de la session.
Exemple de code
Pour des exemples complets, consultez l'exemple de code dans le SDK RocketMQ 5.x gRPC.
Envoyer des messages
Producer producer = provider.newProducerBuilder()
.setTopics(topic)
.setClientConfiguration(clientConfiguration)
.build();
final Message message = provider.newMessageBuilder()
.setTopic(topic)
// Set a message key for precise lookup by keyword.
.setKeys("messageKey")
// Set LiteTopic
.setLiteTopic("lite-topic-1")
// Message body
.setBody("messageBody".getBytes())
.build();
try {
final SendReceipt sendReceipt = producer.send(message);
log.info("Send message successfully, messageId={}", sendReceipt.getMessageId());
} catch (LiteTopicQuotaExceededException e) {
// LiteTopic quota exceeded. Evaluate and increase quota.
log.error("Lite topic quota exceeded", e);
} catch (Throwable t) {
log.error("Failed to send message", t);
}
Consommer des messages
Utilisez la classe LitePushConsumer :
// Initialize LitePushConsumer with consumer group, target topic, and communication parameters.
LitePushConsumer litePushConsumer = provider.newLitePushConsumerBuilder()
.setClientConfiguration(clientConfiguration)
// Topic bound to the ConsumerGroup when created in the console
.bindTopic(topicName)
// Set consumer group
.setConsumerGroup(consumerGroup)
.setMessageListener(messageView -> {
// Process message and return consumption result.
LOGGER.info("Consume message={}", messageView);
return ConsumeResult.SUCCESS;
})
.build();
try {
// Subscribe to desired LiteTopics
litePushConsumer.subscribeLite("lite-topic-1");
litePushConsumer.subscribeLite("lite-topic-2");
litePushConsumer.subscribeLite("lite-topic-3");
} catch (LiteSubscriptionQuotaExceededException e) {
// LiteTopic subscription quota exceeded. Evaluate and increase quota.
log.error("Lite subscription quota exceeded", e);
} catch (Throwable t) {
log.error("Failed to subscribe lite topic", t);
}
// After business processing, unsubscribe from unused LiteTopics promptly
litePushConsumer.unsubscribeLite("lite-topic-3");
// Get current set of subscribed LiteTopics
Set<String> liteTopicSet = litePushConsumer.getLiteTopicSet();
Mettre à jour dynamiquement les abonnements
/**
* Dynamically add a subscription.
* The subscribeLite() method makes a network call and validates quotas,
* so it might fail.
* Always check the result to confirm successful subscription.
* Possible failure scenarios:
* 1. Network error – retry the call.
* 2. Quota validation fails – throws LiteSubscriptionQuotaExceededException.
* Evaluate whether your quota meets requirements, and promptly call
* unsubscribeLite() to release resources for unused topics.
*/
litePushConsumer.subscribeLite("lite-topic-1");
// Dynamically remove a subscription
litePushConsumer.unsubscribeLite("lite-topic-1");
Limites
Un seul consommateur peut s'abonner à un maximum de 2 000 LiteTopics (ajustable via ticket).
Chaque LiteTopic prend en charge un TPS de consommation maximal de 200.
-
Pour garantir la stabilité du service, chaque instance applique une limite sur le nombre total de LiteTopics pouvant être créés ou faisant l'objet d'un abonnement. Consultez le tableau suivant pour connaître les quotas spécifiques (ajustables via ticket).
-
Nombre de LiteTopics
Définition : Nombre total de LiteTopics actuellement créés et actifs sous une seule instance durant son cycle de vie.
Déclenchement et impact : Lorsque cette limite est atteinte, les tentatives d'envoi d'un message vers un LiteTopic inexistant (ce qui déclencherait une création automatique) échouent avec une erreur d'envoi.
-
Nombre d'abonnements LiteTopic
Définition : Nombre total de relations d'abonnement actives entre tous les clients consommateurs en ligne et les LiteTopics sous une instance. Ce nombre change dynamiquement.
Impact : Lorsque cette limite est atteinte, toute tentative d'un consommateur de s'abonner à un nouveau LiteTopic échoue.
Règle spéciale : Même si un LiteTopic est supprimé, tous les abonnements consommateurs restants à celui-ci continuent de compter dans le total jusqu'à ce que ces consommateurs se désabonnent.
-
Instances Serverless
|
Architecture de déploiement |
Mode de capacité |
Spécifications |
Nombre maximal de LiteTopics créables ou faisant l'objet d'un abonnement |
|
Dédié |
Réservé + élastique |
5000 |
300 000 |
|
10000 |
600 000 |
||
|
15000 |
720 000 |
||
|
[20 000, 50 000] |
1 000 000 |
||
|
(50 000, 100 000] |
1 500 000 |
||
|
(100 000, 200 000] |
2 400 000 |
||
|
(200 000, 300 000] |
4 700 000 |
||
|
(300 000, 500 000] |
6 300 000 |
||
|
(500 000, 1 000 000] |
11 600 000 |
Instances non Serverless (abonnement et paiement à l'utilisation)
Édition Standard
|
Type d'instance |
Limite TPS de base pour l'envoi et la réception (ops/sec) |
Nombre maximal de LiteTopics créables ou faisant l'objet d'un abonnement |
|
rmq.s2.2xlarge |
2000 |
150 000 |
|
rmq.s2.4xlarge |
4000 |
250 000 |
|
rmq.s2.6xlarge |
6000 |
300 000 |
Édition Professionnelle
|
Type d'instance |
Limite TPS de base pour l'envoi et la réception (ops/sec) |
Nombre maximal de LiteTopics créables ou faisant l'objet d'un abonnement |
|
rmq.p2.2xlarge |
2000 |
150 000 |
|
rmq.p2.4xlarge |
4000 |
250 000 |
|
rmq.p2.6xlarge |
6000 |
300 000 |
|
rmq.p2.10xlarge |
10000 |
600 000 |
|
rmq.p2.20xlarge |
20000 |
800 000 |
|
rmq.p2.30xlarge |
30000 |
1 000 000 |
|
rmq.p2.40xlarge |
40000 |
1,2 million |
|
rmq.p2.50xlarge |
50000 |
1,4 million |
|
rmq.p2.100xlarge |
100000 |
2 200 000 |
|
rmq.p2.120xlarge |
120000 |
2,7 millions |
|
rmq.p2.150xlarge |
150000 |
3,3 millions |
|
rmq.p2.200xlarge |
200000 |
4,5 millions |
Édition Platinum
|
Type d'instance |
Limite TPS de base pour l'envoi et la réception (ops/sec) |
Nombre maximal de LiteTopics créables ou faisant l'objet d'un abonnement |
|
rmq.u2.10xlarge |
10000 |
600 000 |
|
rmq.u2.20xlarge |
20000 |
800 000 |
|
rmq.u2.30xlarge |
30000 |
1 000 000 |
|
rmq.u2.40xlarge |
40000 |
1 200 000 |
|
rmq.u2.50xlarge |
50000 |
1 400 000 |
|
rmq.u2.60xlarge |
60000 |
1 600 000 |
|
rmq.u2.70xlarge |
70000 |
1 700 000 |
|
rmq.u2.80xlarge |
80000 |
1 800 000 |
|
rmq.u2.90xlarge |
90000 |
2 000 000 |
|
rmq.u2.100xlarge |
100000 |
2 200 000 |
|
rmq.u2.120xlarge |
120000 |
2 700 000 |
|
rmq.u2.150xlarge |
150000 |
3 300 000 |
|
rmq.u2.200xlarge |
200000 |
4 500 000 |
|
rmq.u2.250xlarge |
250000 |
5 600 000 |
|
rmq.u2.300xlarge |
300000 |
6 300 000 |
|
rmq.u2.350xlarge |
350000 |
7 500 000 |
|
rmq.u2.400xlarge |
400000 |
9 300 000 |
|
rmq.u2.450xlarge |
450000 |
10 400 000 |
|
rmq.u2.500xlarge |
500000 |
11 600 000 |
|
rmq.u2.550xlarge |
550000 |
12 800 000 |
|
rmq.u2.600xlarge |
600000 |
14 000 000 |
|
rmq.u2.1000xlarge |
1000000 |
23 200 000 |