ApsaraMQ for RocketMQ utilise les décalages de consommation (consumer offsets) pour gérer la progression des consommateurs. Cette rubrique décrit le mécanisme de gestion de la progression des consommateurs dans ApsaraMQ for RocketMQ.
Contexte
Dans ApsaraMQ for RocketMQ, les messages peuvent être générés avant ou après leur souscription par les consommateurs. Comment un consommateur sait-il où commencer à consommer les messages et comment marquer les messages déjà traités ? Pour répondre à cette problématique, ApsaraMQ for RocketMQ a mis en place un mécanisme de gestion de la progression des consommateurs.
Le mécanisme de gestion de la progression des consommateurs de ApsaraMQ for RocketMQ résout les problèmes suivants :
Où un client commence-t-il à consommer les messages après son démarrage ?
Comment marquer un message consommé pour éviter tout double traitement ?
Un même client peut-il reconsommer un message en cas d'exception de service ?
Mécanisme de fonctionnement
Décalage (Offset)
Dans ApsaraMQ for RocketMQ, les messages sont stockés dans plusieurs files d'attente d'un topic spécifique, selon l'ordre de leur arrivée sur le broker. Chaque message reçoit une coordonnée unique de type Long, également appelée décalage (offset) du message.
Théoriquement, une file d'attente de messages peut stocker un nombre illimité de messages. Par conséquent, la plage de valeurs du décalage s'étend de 0 à Long.MAX_VALUE. Vous pouvez localiser un message en vous basant sur son topic, sa file d'attente et son décalage. La figure suivante illustre la relation entre ces concepts.
Dans ApsaraMQ for RocketMQ, le décalage du premier message d'une file d'attente est appelé décalage minimum (MinOffset), et celui du dernier message est appelé décalage maximum (MaxOffset). Bien qu'une file d'attente de messages puisse théoriquement contenir un nombre illimité de messages, l'espace de stockage des machines physiques qui les hébergent est limité. Par conséquent, ApsaraMQ for RocketMQ supprime dynamiquement les messages les plus anciens d'une file d'attente, ce qui entraîne une augmentation constante des valeurs MinOffset et MaxOffset de cette file.
Décalage du consommateur (Consumer offset)
ApsaraMQ for RocketMQ suit le modèle publication-abonnement. Plusieurs groupes de consommateurs peuvent s'abonner à la même file d'attente. Dans ce type de scénario, si un consommateur supprimait un message après l'avoir consommé, les autres consommateurs ne pourraient plus le traiter.
Pour éviter cette situation, ApsaraMQ for RocketMQ utilise les décalages de consommation afin de gérer la progression de chaque consommateur. ApsaraMQ for RocketMQ ne supprime pas immédiatement un message après sa consommation. Le service conserve plutôt un enregistrement du dernier message consommé par un groupe de consommateurs, appelé décalage du consommateur (consumer offset).
En cas de redémarrage d'un client, le consommateur reprend le traitement des messages en se basant sur le décalage enregistré dans le broker. Si le décalage du consommateur expire et est supprimé, la valeur MinOffset de la file d'attente enregistrée dans le broker sert de décalage de remplacement.
Les décalages de consommation sont enregistrés et restaurés depuis les brokers ApsaraMQ for RocketMQ et ne sont pas liés à un consommateur spécifique. Par conséquent, ApsaraMQ for RocketMQ permet de restaurer la progression des consommateurs sur différents clients.
La figure suivante montre les relations entre le décalage minimum, le décalage maximum et le décalage d'un consommateur dans une file d'attente de messages.
-
Le décalage du consommateur est toujours inférieur ou égal au décalage maximum.
Si les messages sont produits et consommés à la même vitesse et qu'aucun message non consommé n'existe dans la file d'attente, le décalage du consommateur est identique au décalage maximum.
Si la consommation est plus lente que la production, des messages non consommés s'accumulent dans la file d'attente. Dans ce cas, le décalage du consommateur est inférieur au décalage maximum, et la différence correspond au nombre de messages non consommés.
Généralement, le décalage du consommateur est supérieur ou égal au décalage minimum. Si le décalage du consommateur est inférieur au décalage minimum, le consommateur ne peut pas consommer de messages. Dans ce cas, le broker restaure le décalage correct pour le consommateur.
Décalage initial du consommateur
Le décalage initial du consommateur correspond au décalage enregistré dans un broker lorsqu'un groupe de consommateurs commence à consommer une file d'attente pour la première fois.
ApsaraMQ for RocketMQ utilise le décalage maximum d'une file d'attente comme décalage initial lorsque le consommateur obtient des messages de cette file pour la première fois. Autrement dit, le consommateur commence la consommation à partir du dernier message de la file d'attente.
Réinitialiser le décalage d'un consommateur
Si le décalage initial ou actuel du consommateur n'est pas aligné avec l'état de votre activité métier, vous pouvez réinitialiser ce décalage pour ajuster la progression de votre consommateur.
Scénarios
Décalage initial inapproprié : le décalage initial correspond au décalage maximum de la file d'attente, et le client commence la consommation à partir du dernier message. Si vous devez consommer des messages plus anciens, réinitialisez le décalage du consommateur à la position d'un message antérieur.
Retard de consommation : un grand nombre de messages peuvent s'accumuler si le consommateur ne parvient pas à suivre la vitesse de génération des messages. Si les messages accumulés ne sont pas critiques, définissez le décalage du consommateur sur une valeur plus élevée pour ignorer ces messages et alléger la charge en aval.
Rétrotraitement et correction métier : si vous souhaitez reconsommer des messages incorrectement traités en raison d'erreurs métier, définissez le décalage du consommateur sur une valeur inférieure.
Fonctionnalité de réinitialisation du décalage
La fonctionnalité de réinitialisation du décalage de ApsaraMQ for RocketMQ offre les capacités suivantes :
-
Réinitialiser le décalage d'un consommateur au dernier décalage
Les consommateurs du groupe spécifié ignorent tous les messages accumulés dans le topic indiqué et commencent la consommation à partir du dernier décalage.
-
Réinitialiser le décalage d'un consommateur à un moment précis
Les consommateurs commencent la consommation à partir du message correspondant au moment de réinitialisation, que ce message ait déjà été consommé ou non.
Vous pouvez spécifier un moment situé dans la plage comprise entre l'envoi du premier message au topic et l'envoi du dernier message.
Si vous réinitialisez le décalage d'un consommateur à un moment précis, le broker ajuste le décalage à la valeur la plus proche de ce point temporel.
Méthodes de configuration
-
Opérations via la console :
Connectez-vous à la console ApsaraMQ for RocketMQ. Dans le volet de navigation de gauche, cliquez sur Instances.
Sur la page Instances, sélectionnez l'instance que vous souhaitez gérer. Dans le volet de navigation de gauche de la page Instance Details qui s'affiche, cliquez sur Groups.
Sur la page Groups, cliquez sur le groupe que vous souhaitez gérer. Sur la page Group Details qui s'affiche, réinitialisez le décalage du consommateur.
Opération API : ResetConsumeOffset
Limites
Après la réinitialisation du décalage d'un consommateur, celui-ci commence à consommer les messages à partir du nouveau décalage. Dans les scénarios de rétrotraitement, le consommateur commence par des messages historiques qui constituent souvent des données froides. Cela engendre des lectures à froid (cold reads) susceptibles de solliciter inutilement votre système. Évaluez les risques et les avantages avant de réinitialiser un décalage. Nous vous recommandons de mettre en œuvre des politiques de contrôle strictes pour cette autorisation afin d'éviter les abus et les réinitialisations fréquentes.
ApsaraMQ for RocketMQ vous permet de réinitialiser le décalage du consommateur uniquement pour les messages visibles. Vous ne pouvez pas réinitialiser le décalage des messages dont l'état est « scheduling » ou « retry pending ». Pour plus d'informations, consultez les rubriques Messages planifiés et différés et Nouvelle tentative de consommation.
Compatibilité des versions
La définition du décalage initial du consommateur varie selon les versions d'ApsaraMQ for RocketMQ utilisées par les brokers :
Dans les versions 4.x et 3.x, le décalage initial du consommateur est défini comme l'état du message d'une file d'attente.
Dans les versions 5.x, le décalage initial du consommateur correspond au décalage maximum de la file d'attente au moment où le consommateur commence à recevoir les messages.
Par conséquent, lors d'une mise à niveau depuis une version antérieure, portez une attention particulière au décalage initial du consommateur lors du lancement de votre client.
Remarques d'utilisation
Contrôlez strictement les autorisations de réinitialisation
La réinitialisation du décalage d'un consommateur impose une charge supplémentaire au système et peut affecter les opérations de lecture et d'écriture des messages. Par conséquent, nous vous recommandons d'évaluer les risques et les avantages avant d'effectuer cette opération.