Tous les produits
Search
Centre de documentation

ApsaraMQ for Kafka:Best practices for consumers

Dernière mise à jour :Aug 11, 2026

Lorsque les consommateurs ne parviennent pas à suivre le rythme des producteurs, les messages s'accumulent et la latence de traitement augmente brusquement. Une mauvaise configuration des offsets peut entraîner le saut ou le retraitement silencieux de messages. Ce guide présente des modèles éprouvés pour créer des applications de consommation fiables et à haut débit sur ApsaraMQ for Kafka, en couvrant les groupes de consommateurs, la gestion des offsets, la gestion des erreurs et l'optimisation des performances.

Fonctionnement de la consommation de messages

Un consommateur ApsaraMQ for Kafka répète un cycle en trois étapes :

  1. Poll : récupère un lot de messages depuis le broker.

  2. Process : exécute votre logique métier sur les messages récupérés.

  3. Poll again : récupère le lot suivant une fois le traitement terminé.

Consumer polling cycle

Maintenez un temps de traitement court pour conserver un intervalle de poll régulier. Des délais de traitement prolongés déclenchent des rééquilibrages ou provoquent une accumulation de messages.

Groupes de consommateurs et équilibrage de charge

Un groupe de consommateurs est un ensemble d'instances de consommateurs qui partagent le même group.id. ApsaraMQ for Kafka distribue les partitions des topics abonnés uniformément entre les consommateurs d'un groupe, de sorte que chaque message est remis à exactement un consommateur.

Exemple : Le groupe A s'abonne au topic A. Trois consommateurs (C1, C2 et C3) font partie du groupe. Chaque message entrant est envoyé à l'un des consommateurs C1, C2 ou C3, ce qui permet d'équilibrer la charge de consommation entre les trois.

ApsaraMQ for Kafka déclenche un rééquilibrage lorsque les consommateurs sont :

  • Démarrés pour la première fois

  • Redémarrés

  • Ajoutés au groupe ou supprimés du groupe

Éviter les rééquilibrages fréquents

Les rééquilibrages suspendent la consommation pendant la réaffectation des partitions. Des rééquilibrages fréquents perturbent le débit et augmentent la latence de traitement. Ils se produisent lorsque les signaux de maintien de connexion (heartbeats) des consommateurs expirent.

Pour éviter les rééquilibrages fréquents, ajustez les paramètres pertinents ou augmentez la vitesse de consommation. Pour des étapes de dépannage détaillées, consultez Pourquoi des rééquilibrages se produisent-ils fréquemment sur mon client consommateur ?

Configuration des partitions

Le nombre de partitions contrôle le nombre de consommateurs pouvant traiter les messages en parallèle. Chaque partition est attribuée à exactement un consommateur au sein d'un groupe ; si vous disposez de plus de consommateurs que de partitions, les consommateurs supplémentaires restent inactifs.

Nombre de partitions recommandé

Le nombre de partitions par défaut dans la console ApsaraMQ for Kafka est de 12, ce qui convient à la plupart des charges de travail. Lors de la mise à l'échelle, suivez ces recommandations :

Nombre de partitions Impact
Inférieur à 12 Peut réduire les performances de production et de consommation de messages
12 -- 100 Plage recommandée pour la plupart des charges de travail
Supérieur à 100 Peut déclencher des rééquilibrages fréquents des consommateurs
Important

Vous ne pouvez pas réduire le nombre de partitions après une augmentation. Augmentez les partitions de manière incrémentielle plutôt que par grands sauts.

Modèles d'abonnement

ApsaraMQ for Kafka prend en charge deux modèles d'abonnement.

Un groupe s'abonne à plusieurs topics

Un seul groupe de consommateurs peut s'abonner à plusieurs topics. Les messages de tous les topics abonnés sont distribués uniformément entre les consommateurs du groupe.

String topicStr = kafkaProperties.getProperty("topic");
String[] topics = topicStr.split(",");
for (String topic: topics) {
    subscribedTopics.add(topic.trim());
}
consumer.subscribe(subscribedTopics);

Plusieurs groupes s'abonnent à un topic

Plusieurs groupes de consommateurs peuvent s'abonner indépendamment au même topic. Chaque groupe reçoit tous les messages, et les groupes fonctionnent indépendamment sans s'affecter mutuellement.

Utilisez ce modèle lorsque différentes applications doivent traiter le même flux de données de manière indépendante, par exemple un groupe pour l'analytique en temps réel et un autre pour l'archivage des données.

Conserver un groupe par application

Utilisez un groupe de consommateurs dédié pour chaque application. Si une seule application doit exécuter différentes logiques de consommation, créez des fichiers de configuration distincts (par exemple, kafka1.properties et kafka2.properties) avec des valeurs group.id distinctes.

Faire correspondre les abonnements au sein d'un groupe

Tous les consommateurs d'un même groupe doivent s'abonner au même ensemble de topics. Des abonnements non concordants compliquent le dépannage et peuvent entraîner un comportement de rééquilibrage inattendu.

Gérer les offsets des consommateurs

Chaque partition suit un offset maximum, c'est-à-dire le nombre total de messages reçus. Chaque consommateur suit un offset de consommateur, c'est-à-dire le nombre de messages qu'il a traités dans cette partition. La différence entre les deux constitue l'accumulation de messages (backlog non consommé).

Validation automatique vs validation manuelle

ApsaraMQ for Kafka fournit deux paramètres pour valider les offsets des consommateurs :

Paramètre Par défaut Description
enable.auto.commit true Active la validation automatique des offsets
auto.commit.interval.ms 1000 (ms) Intervalle entre les validations automatiques

Lorsque la validation automatique est activée, le client vérifie le temps écoulé depuis la dernière validation avant chaque poll. Si le temps écoulé dépasse auto.commit.interval.ms, il valide l'offset actuel.

Important

Lorsque vous utilisez la validation automatique, assurez-vous que tous les messages du poll précédent sont entièrement traités avant le poll suivant. Si le consommateur écrit dans un magasin de données externe et que cette écriture échoue, la validation automatique peut tout de même avancer l'offset, ce qui entraîne le saut définitif de ces messages au lieu de leur nouvelle tentative. Comprenez ce compromis avant de vous fier au comportement par défaut.

Valider les offsets manuellement

Pour valider les offsets manuellement, définissez enable.auto.commit sur false et appelez la fonction commit(offsets) après le traitement de chaque lot.

Comportement de réinitialisation des offsets

Les offsets des consommateurs sont réinitialisés dans deux situations :

  • Aucun offset validé n'existe, par exemple lorsqu'un consommateur se connecte à un broker pour la première fois.

  • L'offset validé n'est pas valide, par exemple si l'offset maximum d'une partition est 10, mais que le consommateur tente de lire à partir de l'offset 11.

Contrôlez le comportement de réinitialisation sur les clients Java avec auto.offset.reset :

Valeur Comportement
latest Démarrer à partir du message le plus récent. Utilisez ceci comme valeur par défaut
earliest Démarrer à partir du message disponible le plus ancien
none Lever une exception au lieu de réinitialiser
Remarque

Privilégiez latest plutôt que earliest pour éviter de retraiter l'historique complet du topic lorsqu'un offset non valide se produit. Si vous validez les offsets manuellement avec la fonction commit(offsets), none est un choix sûr car vos offsets devraient toujours être valides.

Gérer les messages volumineux

Lorsque des messages individuels dépassent les tailles habituelles, ajustez ces paramètres pour contrôler le comportement de récupération :

Paramètre Recommandation
max.poll.records Nombre maximal d'enregistrements par poll. Définissez sur 1 pour les messages supérieurs à 1 Mo
fetch.max.bytes Données maximales par requête de récupération. Définissez une valeur légèrement supérieure à la taille attendue du message
max.partition.fetch.bytes Données maximales par partition et par récupération. Définissez une valeur légèrement supérieure à la taille attendue du message

Avec ces paramètres, le consommateur récupère les messages volumineux un par un, évitant ainsi la pression mémoire due à des lots de taille excessive.

Gérer la duplication des messages

ApsaraMQ for Kafka utilise la sémantique de livraison au moins une fois : chaque message est livré au moins une fois pour éviter la perte de données, mais des duplications peuvent survenir lors d'erreurs réseau ou de redémarrages du client.

Implémenter une consommation idempotente

Si votre application est sensible aux duplications, par exemple pour les transactions financières ou le traitement des commandes, implémentez la déduplication au niveau de l'application :

  1. Joignez une clé unique à chaque message lors de la production (par exemple, un ID de transaction).

  2. Avant de traiter un message consommé, vérifiez si la clé a déjà été traitée.

  3. Ignorez les messages dont les clés ont déjà été traitées.

Si votre application tolère des duplications occasionnelles (par exemple, l'agrégation de métriques ou l'ingestion de journaux), ignorez cette étape.

Gérer les échecs de consommation

Les messages au sein d'une partition sont consommés séquentiellement. En cas d'échec du traitement, par exemple dû à des données mal formées ou à une panne d'un service en aval, choisissez l'une de ces stratégies :

Stratégie Compromis
Nouvelle tentative sur place Retentez le message ayant échoué jusqu'à ce qu'il réussisse. Simple à mettre en œuvre, mais bloque le thread du consommateur et provoque une accumulation de messages sur cette partition
Topic de lettres mortes Envoyez les messages ayant échoué vers un topic dédié pour inspection ultérieure. La consommation continue sans blocage, mais nécessite un processus distinct pour enquêter et retraiter les échecs

Pour la plupart des charges de travail en production, l'approche par lettres mortes évite le blocage du pipeline de consommation.

Réduire la latence de consommation

ApsaraMQ for Kafka utilise un modèle basé sur la récupération (pull) : les consommateurs récupèrent les messages du broker à chaque cycle de poll. Lorsque les consommateurs suivent le rythme des producteurs, la latence reste faible. Une augmentation soudaine de la latence indique généralement une accumulation de messages.

Causes courantes d'accumulation

Cause Solution
La vitesse de consommation est plus lente que la vitesse de production Ajoutez des consommateurs ou augmentez le débit de traitement
Le thread du consommateur est bloqué par un appel distant lent Définissez des délais d'expiration sur les appels externes pour libérer le thread après une attente limitée

Augmenter le débit de consommation

Ajoutez plus de consommateurs. Démarrez des instances de consommateurs supplémentaires (un thread par consommateur) ou déployez davantage de processus de consommateurs. Le nombre de consommateurs actifs ne peut pas dépasser le nombre de partitions ; les supplémentaires restent inactifs.

Utilisez un pool de threads pour le traitement. Découplez la récupération du traitement pour paralléliser le travail :

  1. Définissez un pool de threads avec un nombre fixe de threads de travail.

  2. Récupérez un lot de messages.

  3. Soumettez chaque message au pool de threads pour un traitement concurrent.

  4. Une fois toutes les tâches du lot terminées, récupérez le lot suivant.

Ce modèle maintient la boucle de poll réactive tout en distribuant le travail de traitement sur plusieurs threads.

Filtrer les messages

ApsaraMQ for Kafka ne propose pas de filtrage de messages intégré. Implémentez le filtrage au niveau de l'application en utilisant l'une de ces approches :

Approche Quand l'utiliser
Topics séparés Un petit nombre de catégories de messages distinctes
Filtrage côté client De nombreuses catégories, ou une logique de filtrage qui change fréquemment

Combinez les deux approches lorsque votre cas d'utilisation l'exige : acheminez les grandes catégories vers différents topics, puis appliquez des filtres granulaires côté consommateur.

Diffuser des messages

ApsaraMQ for Kafka ne prend pas en charge nativement la diffusion de messages. Pour remettre chaque message à plusieurs consommateurs indépendants, créez un groupe de consommateurs distinct pour chaque consommateur devant recevoir tous les messages. Chaque groupe consomme indépendamment le flux complet de messages du topic.