Tous les produits
Search
Centre de documentation

ApsaraMQ for Kafka:Best practices for producers

Dernière mise à jour :Aug 10, 2026

Lorsqu'un producteur Kafka envoie un grand volume de messages, une configuration incorrecte des nouvelles tentatives, du regroupement par lots ou des accusés de réception peut entraîner une perte de données, une latence excessive ou des erreurs de mémoire insuffisante (out-of-memory). Cette rubrique explique comment configurer votre producteur ApsaraMQ for Kafka pour garantir la fiabilité des envois et optimiser le débit. Tous les exemples utilisent le client Java. Les concepts fondamentaux s'appliquent à d'autres langages, mais les noms des paramètres et les détails d'implémentation peuvent différer.

Envoyer un message

Toute interaction avec un producteur commence par producer.send(), qui accepte un ProducerRecord contenant le topic, la partition, l'horodatage, la clé et la valeur.

Future<RecordMetadata> metadataFuture = producer.send(new ProducerRecord<String, String>(
        topic,   // The message topic.
        null,   // The partition number. Set this to null to let the producer assign a partition.
        System.currentTimeMillis(),   // The timestamp.
        String.valueOf(value.hashCode()),   // The message key.
        value   // The message value.
));

La méthode send() est asynchrone. Pour bloquer l'exécution jusqu'à l'obtention du résultat, appelez :

RecordMetadata metadata = metadataFuture.get(timeout, TimeUnit.MILLISECONDS);

Pour des exemples complets de SDK, consultez Présentation du SDK.

Clé et valeur

Les messages dans ApsaraMQ for Kafka 0.10.2 comportent deux champs :

  • Key : identifiant du message. Définissez une clé unique par message afin de suivre son cycle de vie via les journaux d'envoi et de consommation.

  • Value : contenu du message.

Pour les envois à haut volume ne nécessitant pas de routage basé sur une clé, omettez la clé et utilisez la stratégie de partitionnement sticky. Cela améliore l'efficacité du regroupement par lots.

Important

ApsaraMQ for Kafka 0.11.0 et versions ultérieures prennent en charge les en-têtes. Pour utiliser les en-têtes, mettez à niveau le serveur vers la version 2.2.0.

Sécurité des threads

Une instance de producteur est thread-safe et peut envoyer des messages vers n'importe quel topic. Utilisez un seul producteur par application.

Configurer les nouvelles tentatives

Dans un environnement distribué, les envois peuvent échouer en raison de problèmes réseau. L'échec peut se produire à deux niveaux : le message a atteint le broker, mais l'accusé de réception (ACK) a été perdu, ou le message n'a jamais atteint le broker.

ApsaraMQ for Kafka utilise une architecture réseau d'adresse IP virtuelle (VIP) qui ferme automatiquement les connexions inactives. Les clients inactifs reçoivent souvent une erreur connection reset by peer. Les nouvelles tentatives permettent de gérer ces échecs temporaires.

Paramètre Description Valeur recommandée
retries Nombre maximal de tentatives pour un envoi ayant échoué. Utilisez la valeur par défaut (définie par la version du client).
retry.backoff.ms Délai entre les nouvelles tentatives, en millisecondes. 1000

Accusés de réception

Le paramètre acks contrôle le nombre de réplicas qui doivent confirmer une écriture avant que le broker ne réponde. Choisissez un paramétrage en fonction de vos exigences en matière de durabilité et de performance.

Paramètre Comportement Débit Durabilité
acks=0 Aucune réponse du broker requise. Le plus élevé La plus faible : perte de données probable en cas de panne
acks=1 Réponse après l'écriture des données par le nœud principal. Moyen Moyenne : perte de données si le nœud principal tombe en panne avant la réplication
acks=all Réponse après l'écriture des données par le nœud principal et tous les nœuds réplica synchronisés (ISR). Le plus faible La plus élevée : perte de données uniquement si le nœud principal et les nœuds réplica tombent en panne simultanément

Pour améliorer les performances d'envoi, définissez acks=1.

Optimiser le regroupement par lots

Le producteur regroupe les messages destinés à la même partition en lots avant l'envoi. Des lots plus grands réduisent le nombre de requêtes réseau, diminuent l'utilisation du CPU et améliorent le débit et la latence. De petits lots provoquent une mise en file d'attente des requêtes côté client et serveur.

Deux paramètres contrôlent le regroupement par lots :

Paramètre Description Valeur par défaut Recommandation
batch.size Taille maximale du lot par partition, en octets. Une requête réseau est déclenchée lorsque le lot atteint cette taille. 16384 (16 Ko) Conservez la valeur par défaut de 16384. Une valeur trop faible dégrade les performances et la stabilité.
linger.ms Temps maximal pendant lequel un message attend dans le tampon avant d'être envoyé, en millisecondes. Lorsque ce délai expire, le producteur envoie le lot indépendamment de la valeur de batch.size. 0 Définissez une valeur comprise entre 100 et 1000.

Un lot est envoyé dès que l'un ou l'autre des seuils est atteint, selon le premier événement survenant. Pour obtenir un bon équilibre entre débit et latence, définissez batch.size=16384 et linger.ms=1000.

Utilisez un client de version 2,4 ou ultérieure, qui active par défaut la stratégie de partitionnement sticky afin de réduire davantage les envois fragmentés.

Stratégie de partitionnement sticky

Seuls les messages envoyés vers la même partition sont regroupés dans un lot ; la stratégie de partitionnement affecte donc directement l'efficacité du regroupement par lots.

Messages avec clé

Pour les messages dotés d'une clé, le producteur hache la clé et sélectionne une partition en fonction du résultat du hachage. Les messages possédant la même clé sont toujours dirigés vers la même partition.

Messages sans clé

Avant Kafka 2,4, la stratégie par défaut pour les messages sans clé était le round-robin : chaque message était envoyé à la partition suivante dans la séquence. Cette approche dispersait les messages sur toutes les partitions, générant de nombreux petits lots et augmentant la latence.

Kafka 2,4 a introduit la stratégie de partitionnement sticky (KIP-480) pour résoudre ce problème. Au lieu de faire tourner les partitions pour chaque message, le producteur reste sur une partition jusqu'à ce que le lot actuel soit plein, puis sélectionne aléatoirement une nouvelle partition. Avec le temps, les messages sont toujours répartis uniformément sur toutes les partitions, mais les lots sont beaucoup plus grands. Cette approche évite le déséquilibre de partition des messages, réduit la latence et améliore les performances globales du service.

Activer le partitionnement sticky

  • Version du client 2,4 ou ultérieure : le partitionnement sticky est activé par défaut. Aucune configuration n'est nécessaire.

  • Version du client antérieure à 2,4 : implémentez un partitionneur personnalisé et définissez-le via partitioner.class. L'exemple suivant change de partition à un intervalle de temps configurable :

public class MyStickyPartitioner implements Partitioner {

    // Records the time of the last partition switch.
    private long lastPartitionChangeTimeMillis = 0L;
    // Records the current partition.
    private int currentPartition = -1;
    // The partition switch interval. Set the interval as needed.
    private long partitionChangeTimeGap = 100L;

    public void configure(Map<String, ?> configs) {}

    /**
     * Compute the partition for the given record.
     *
     * @param topic The topic name
     * @param key The key to partition on (or null if no key)
     * @param keyBytes serialized key to partition on (or null if no key)
     * @param value The value to partition on or null
     * @param valueBytes serialized value to partition on or null
     * @param cluster The current cluster metadata
     */
    public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {

        // Get all partition information.
        List<PartitionInfo> partitions = cluster.partitionsForTopic(topic);
        int numPartitions = partitions.size();

        if (keyBytes == null) {
            List<PartitionInfo> availablePartitions = cluster.availablePartitionsForTopic(topic);
            int availablePartitionSize = availablePartitions.size();

            // Check the current active partitions.
            if (availablePartitionSize > 0) {
                handlePartitionChange(availablePartitionSize);
                return availablePartitions.get(currentPartition).partition();
            } else {
                handlePartitionChange(numPartitions);
                return currentPartition;
            }
        } else {
            // For a message with a key, select a partition based on the key's hash value.
            return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;
        }
    }

    private void handlePartitionChange(int partitionNum) {
        long currentTimeMillis = System.currentTimeMillis();

        // If the time since the last switch exceeds the switch interval, switch to the next partition. Otherwise, use the current partition.
        if (currentTimeMillis - lastPartitionChangeTimeMillis >= partitionChangeTimeGap
            || currentPartition < 0 || currentPartition >= partitionNum) {
            lastPartitionChangeTimeMillis = currentTimeMillis;
            currentPartition = Utils.toPositive(ThreadLocalRandom.current().nextInt()) % partitionNum;
        }
    }

    public void close() {}

}

Prévenir les erreurs de mémoire insuffisante (OOM)

Le producteur met en cache les messages en mémoire avant d'envoyer les lots. Si le cache devient trop volumineux, une erreur de mémoire insuffisante (OOM) se produit.

Paramètre Description Valeur par défaut Recommandation
buffer.memory Mémoire totale disponible pour le tampon d'envoi du producteur, en octets. Si le tampon est trop petit, l'allocation de mémoire peut prendre beaucoup de temps, ce qui affecte les performances d'envoi et peut provoquer des délais d'expiration d'envoi. 33554432 (32 Mo) Définissez au minimum batch.size x nombre de partitions x 2. La valeur par défaut de 32 Mo est suffisante pour un seul producteur.
Important

L'exécution de plusieurs producteurs dans la même machine virtuelle Java (JVM) multiplie l'utilisation de la mémoire. Chaque producteur alloue son propre buffer.memory : quatre producteurs avec la valeur par défaut de 32 Mo chacun consomment 128 Mo de tas. En production, un seul producteur par application suffit généralement. Si vous devez exécuter plusieurs producteurs, réduisez la valeur de buffer.memory pour chacun ou augmentez la taille du tas JVM en conséquence.

Ordre des partitions

Au sein d'une seule partition, les messages sont stockés et consommés dans l'ordre où ils ont été envoyés.

Par défaut, ApsaraMQ for Kafka ne garantit pas un ordre strict au sein d'une partition. Lors d'une mise à niveau ou d'un basculement, un petit nombre de messages peut être désordonné lorsqu'ils sont redirigés vers une autre partition. Ce compromis améliore la disponibilité.

Pour imposer un ordre strict au sein d'une partition, sélectionnez local storage lors de la création du topic.