Cette rubrique explique comment configurer les paramètres du client pour ApsaraMQ for Kafka. Une configuration adéquate influence directement le débit des messages, la fiabilité de la livraison et la stabilité des consommateurs. Les sections suivantes détaillent les paramètres des producteurs et des consommateurs, en proposant des valeurs recommandées et des conseils d'optimisation pour les charges de travail en production.
Paramètres du producteur
Livraison des messages
acks
Détermine le nombre d'accusés de réception que le producteur exige du broker avant de considérer l'envoi comme réussi.
| Valeur | Comportement | Compromis |
|---|---|---|
0 |
Aucun accusé de réception du broker. | Débit maximal, risque élevé de perte de données. |
1 |
Accusé de réception après l'écriture des données par le leader. | Équilibre entre débit et durabilité. Une perte de données est possible si le leader tombe en panne avant que les followers n'aient répliqué les données. |
all |
Accusé de réception après l'écriture des données par le leader et tous les réplicas synchronisés (in-sync). | Débit minimal, durabilité maximale. La perte de données ne survient que si le leader et tous les réplicas synchronisés tombent en panne simultanément. |
Valeur recommandée : 1 pour la plupart des charges de travail qui privilégient le débit à une durabilité stricte.
retries
Nombre maximal de tentatives de renvoi par le producteur en cas d'échec d'envoi. Une valeur plus élevée aide le producteur à récupérer après des pannes temporaires du broker, telles que les élections de leader. Combinez ce paramètre avec retry.backoff.ms pour contrôler le rythme des nouvelles tentatives.
retry.backoff.ms
Délai entre les tentatives d'envoi. Une valeur trop faible peut provoquer des tempêtes de nouvelles tentatives lors des basculements de broker.
| Recommandé | Par défaut | Unité |
|---|---|---|
| 1000 | -- | millisecondes |
Mise en lots et débit
La mise en lots amortit la surcharge réseau en combinant plusieurs enregistrements dans une seule requête. Deux paramètres contrôlent l'envoi d'un lot : sa taille et le temps d'attente.
batch.size
Taille maximale d'un lot par partition. Lorsqu'un lot atteint cette taille, le producteur l'envoie immédiatement.
| Type | Par défaut | Valeurs valides | Unité |
|---|---|---|---|
| int | 16384 | [0,...] | octets |
Conservez la valeur par défaut de 16384 pour la plupart des charges de travail. Une valeur plus petite augmente le nombre de requêtes réseau et réduit le débit. Si vous augmentez batch.size, assurez-vous que buffer.memory est suffisamment grand pour accueillir les lots plus volumineux.
linger.ms
Temps maximal pendant lequel le producteur attend qu'un lot se remplisse avant de l'envoyer. Ce mécanisme fonctionne comme l'algorithme de Nagle dans TCP : dès qu'un lot atteint batch.size, il est envoyé immédiatement, indépendamment du minuteur d'attente. Si le lot est encore inférieur à batch.size lorsque linger.ms expire, le producteur envoie les données accumulées.
| Recommandé | Par défaut | Unité |
|---|---|---|
| 100 à 1000 | 0 | millisecondes |
Une valeur linger.ms plus élevée améliore l'efficacité de la mise en lots et le débit, au détriment de la latence par message.
Gestion de la mémoire
buffer.memory
Mémoire totale allouée par le producteur pour la mise en tampon des enregistrements non envoyés sur toutes les partitions. Si ce pool est épuisé, send() bloque ou lève une exception selon la valeur de max.block.ms. Un tampon sous-dimensionné entraîne une allocation mémoire lente, une réduction du débit ou des délais d'expiration d'envoi.
Unité : octets. Par défaut : 33554432 (32 Mo).
Formule de dimensionnement :
buffer.memory >= batch.size x number_of_partitions x 2
Par exemple, avec batch.size=16384 et 50 partitions :
16384 x 50 x 2 = 1,638,400 bytes (~1.6 MB minimum)
Si vous augmentez batch.size pour améliorer le débit, augmentez buffer.memory proportionnellement.
Partitionnement
partitioner.class
Détermine comment le producteur attribue les enregistrements aux partitions. La stratégie de partitionnement « sticky » (adhésif) réduit le nombre de lots incomplets en remplissant le lot d'une partition avant de passer à la suivante.
| Version du client Kafka | Stratégie par défaut |
|---|---|
| 2.4 et ultérieure | Partitionnement sticky (par défaut) |
| Antérieure à 2.4 | Round-robin |
Si votre client producteur est antérieur à la version 2.4, définissez explicitement le partitionneur sticky pour améliorer l'efficacité de la mise en lots.
Paramètres du consommateur
Optimisation de la récupération (Fetch)
Ces paramètres contrôlent la quantité de données récupérée par le consommateur lors de chaque requête. Leur optimisation affecte à la fois le débit et la latence.
fetch.min.bytes
Quantité minimale de données que le broker accumule avant de renvoyer une réponse de récupération. Une valeur plus grande réduit la fréquence des récupérations et la charge CPU du broker, ce qui améliore le débit mais augmente la latence de bout en bout des messages. Unité : octets.
Évaluez le taux de messages de votre producteur avant de définir cette valeur. Si les messages arrivent lentement, une valeur fetch.min.bytes élevée ajoute un délai inutile.
fetch.max.wait.ms
Temps maximal pendant lequel le broker attend d'accumuler fetch.min.bytes avant de renvoyer une réponse. Unité : millisecondes.
Le comportement varie selon le type de stockage :
Stockage local : Le broker attend jusqu'à ce que
fetch.min.bytessoit atteint ou quefetch.max.wait.msexpire, selon la première éventualité.Stockage cloud : Le broker renvoie une réponse immédiatement lorsque de nouvelles données arrivent, indépendamment de
fetch.min.bytes.
max.partition.fetch.bytes
Quantité maximale de données renvoyées par le broker par partition lors d'une seule récupération. Unité : octets.
Gestion de session et rééquilibrage
Des paramètres de session et d'interrogation mal configurés sont la cause la plus fréquente de rééquilibrages inattendus des consommateurs. Un rééquilibrage suspend toute la consommation dans le groupe jusqu'à ce que les partitions soient réattribuées ; évitez donc de déclencher des rééquilibrages inutiles.
session.timeout.ms
Temps maximal entre les battements de cœur (heartbeats) avant que le broker ne considère le consommateur comme mort et ne déclenche un rééquilibrage.
| Recommandé | Plage valide | Par défaut | Unité |
|---|---|---|---|
| 30000 à 60000 | 6000 à 300000 | 10000 | millisecondes |
Dans le client Java 0.10.1 et versions ultérieures, un thread d'arrière-plan dédié envoie les battements de cœur indépendamment de
poll()
. Dans les versions Java antérieures ou les clients non Java, les battements de cœur sont envoyés lors des appels
poll()
, donc
session.timeout.ms
doit prendre en compte à la fois le temps de traitement des données et l'intervalle des battements de cœur.
Conseil :
Définissez
heartbeat.interval.ms
à une valeur inférieure ou égale au tiers de
session.timeout.ms
. Par exemple, si
session.timeout.ms
est de 45000, définissez
heartbeat.interval.ms
à 15000 ou moins.
max.poll.records
Nombre maximal d'enregistrements renvoyés lors d'un seul appel poll(). Si le consommateur ne peut pas traiter autant d'enregistrements avant l'échéance du prochain poll(), le broker le considère comme mort et déclenche un rééquilibrage.
Formule de dimensionnement :
max.poll.records < messages_per_thread_per_second x consumer_threads x session_timeout_seconds
Par exemple, avec 500 msg/s par thread, 4 threads et un délai d'expiration de session de 45 secondes :
500 x 4 x 45 = 90,000
Définissez max.poll.records en dessous de cette valeur pour garantir que le consommateur termine toujours le traitement avant l'expiration de la session.
max.poll.interval.ms
Intervalle maximal entre les appels consécutifs poll() avant que le broker ne retire le consommateur du groupe. Ce paramètre s'applique uniquement au client Java 0.10.1 et versions ultérieures, où les battements de cœur s'exécutent sur un thread séparé.
| Par défaut | Unité |
|---|---|
| 300000 | millisecondes |
Formule de dimensionnement :
max.poll.interval.ms > time_per_record x max.poll.records
Dans la plupart des cas, la valeur par défaut de 300000 (5 minutes) est suffisante. Augmentez-la uniquement si votre logique de traitement est exceptionnellement lente.
Gestion des offsets
enable.auto.commit
Contrôle si le consommateur valide automatiquement les offsets à intervalle fixe.
| Valeur | Comportement |
|---|---|
true (par défaut) |
Les offsets sont validés automatiquement toutes les auto.commit.interval.ms millisecondes. Plus simple à gérer, mais peut entraîner un double traitement après un crash. |
false |
Votre application doit appeler explicitement commitSync() ou commitAsync(). Utilisez cette option pour une sémantique « au moins une fois » ou « exactement une fois » lorsqu'elle est combinée à un traitement idempotent. |
auto.commit.interval.ms
Intervalle pour les validations automatiques d'offsets lorsque enable.auto.commit est défini sur true.
| Par défaut | Unité |
|---|---|
| 1000 | millisecondes |
Un intervalle plus court réduit la fenêtre de duplication des messages après un crash, mais augmente le nombre de requêtes de validation envoyées au broker.
auto.offset.reset
Détermine ce qui se produit lorsque le consommateur n'a aucun offset validé ou que l'offset validé n'est plus valide (par exemple, l'offset a été supprimé en raison des politiques de rétention).
| Valeur | Comportement |
|---|---|
latest |
Commence la consommation à partir de l'offset le plus récent. Nouveaux messages uniquement. |
earliest |
Commence la consommation à partir de l'offset disponible le plus ancien. Retraite tous les messages conservés. |
none |
Lève une exception. Utilisez cette option lorsque votre application gère les offsets manuellement. |
Valeur recommandée : latest. L'utilisation de earliest amène le consommateur à retraiter tous les messages conservés chaque fois qu'un offset invalide est rencontré, ce qui peut entraîner un double traitement et une augmentation soudaine du retard du consommateur.
Si votre application gère les offsets manuellement (par exemple, en stockant les offsets dans une base de données externe), définissez ce paramètre sur none.
Référence rapide des paramètres
Paramètres du producteur
| Paramètre | Par défaut | Recommandé | Unité |
|---|---|---|---|
acks |
-- | 1 (débit) ou all (durabilité) |
-- |
retries |
-- | Définir selon les exigences de disponibilité | -- |
retry.backoff.ms |
100 | 1000 | ms |
batch.size |
16384 | 16384 (par défaut) | octets |
linger.ms |
0 | 100 à 1000 | ms |
buffer.memory |
33554432 | >= batch.size x partitions x 2 |
octets |
partitioner.class |
Sticky (2.4+) | Partitionnement sticky | -- |
Paramètres du consommateur
| Paramètre | Par défaut | Recommandé | Unité |
|---|---|---|---|
fetch.min.bytes |
1 | Ajuster en fonction du taux de messages du producteur | octets |
fetch.max.wait.ms |
500 | -- | ms |
max.partition.fetch.bytes |
1048576 | -- | octets |
session.timeout.ms |
10000 | 30000 à 60000 | ms |
heartbeat.interval.ms |
3000 | <= 1/3 de session.timeout.ms |
ms |
max.poll.records |
500 | Voir la formule de dimensionnement | -- |
max.poll.interval.ms |
300000 | 300000 (par défaut) | ms |
enable.auto.commit |
true | Dépend de la sémantique de livraison | -- |
auto.commit.interval.ms |
1000 | 1000 (par défaut) | ms |
auto.offset.reset |
latest |
latest |
-- |