Production de messages
Un producteur est une application qui publie des flux de messages dans des rubriques Kafka. Ces informations se concentrent sur l'interface de programmation Java qui fait partie du projet Apache Kafka. Les concepts s'appliquent également à d'autres langages, mais avec des noms parfois légèrement différents.
Dans les interfaces de programmation, un message est appelé un enregistrement. Ainsi, la classe Java org.apache.kafka.clients.producer.ProducerRecord est utilisée pour représenter un message du point de vue de l'API producteur. Les termes enregistrement et message peuvent être utilisés de manière interchangeable, mais, en résumé, un enregistrement est utilisé pour représenter un message.
Lorsqu'un producteur se connecte à Kafka, il établit une connexion d'amorçage initiale. Cette connexion peut viser n'importe quel serveur du cluster. Le producteur demande des informations concernant la partition et le responsable de la rubrique dans laquelle il veut publier. Ensuite, le producteur établit une autre connexion avec le chef de partition et peut publier des messages. Ces actions s'exécutent automatiquement en interne lorsque votre producteur se connecte au cluster Kafka.
Pour garantir la disponibilité, les courtiers Kafka répliquent les messages, de sorte que si un courtier n'est pas disponible, les autres peuvent toujours recevoir les messages des producteurs et les envoyer aux consommateurs. Event Streams utilise un facteur de réplication de 3, ce qui signifie que chaque message est stocké sur trois courtiers. Un message envoyé au responsable de la partition n'est pas immédiatement disponible pour les consommateurs. Le responsable ajoute l'enregistrement du message à la partition, en lui affectant le numéro de position suivant de cette partition. Une fois que tous les suiveurs des répliques synchronisées ont répliqué l'enregistrement et confirmé qu'ils ont écrit l'enregistrement dans leurs répliques, l'enregistrement est maintenant validé et devient disponible pour les consommateurs.
Chaque message est représenté par un enregistrement composé de deux parties : la clé et la valeur. La clé est généralement utilisée pour les données concernant le message et la valeur constitue le corps du message. Étant donné que de nombreux outils de l'écosystème Kafka (tels que les connecteurs vers d'autres systèmes) n'utilisent que la valeur et ignorent la clé, il est préférable de placer toutes les données du message dans la valeur et d'utiliser la clé pour le partitionnement ou le compactage du journal. Il ne faut pas s'attendre à ce que tout ce qui lit dans Kafka utilise la clé.
Bon nombre d'autres systèmes de messagerie disposent également d'un moyen de transmettre d'autres informations avec les messages. La version 0.11 de Kafka introduit des en-têtes d'enregistrement à cette fin.
Il peut s'avérer utile de lire ces informations avec les messages de consommation dans Event Streams.
Paramètres de configuration
De nombreux paramètres de configuration existent pour le fournisseur. Vous pouvez contrôler certains aspects du producteur, notamment la mise en lots, les nouvelles tentatives et l'accusé de réception des messages. Les plus importants sont décrits ci après :
| Nom | Description | Valeur valides | Valeur par défaut |
|---|---|---|---|
| key.serializer | Classe utilisée pour sérialiser les clés. | Classe Java qui implémente l'interface Serializer, telle que org.apache.kafka.common.serialization.StringSerializer. | Pas de valeur par défaut - vous devez spécifier une valeur. |
| value.serializer | Classe utilisée pour sérialiser les valeurs. | Classe Java qui implémente l'interface Serializer, telle que org.apache.kafka.common.serialization.StringSerializer. | Pas de valeur par défaut - vous devez spécifier une valeur. |
| acks | Nombre de serveurs requis pour accuser réception de chaque message publié. Ce paramètre contrôle les garanties de durabilité requises par le producteur. | 0, 1, all (ou -1) | all (Kafka 3.0 et versions ultérieures) 1 (antérieur à Kafka 3.0) |
| nouveaux essais | Nombre de fois où le client renvoie un message lorsque l'envoi rencontre une erreur. | 0,... | 0 |
| max.block.ms | Nombre de millisecondes pendant lesquelles un envoi ou une demande de métadonnées peut rester bloqué en attente. | 0,... | 60000 (1 minute) |
| max.in.flight.requests.per.connection | Nombre maximal de demandes non acquittées que le client envoie sur une connexion avant de bloquer les autres demandes. | 1,... | 5 |
| request.timeout.ms | Durée maximale d'attente du producteur pour une réponse à une demande. Si la réponse n'est pas reçue avant l'expiration du délai, la demande est relancée ou échoue si le nombre de tentatives est épuisé. | 0,... | 30000 (30 secondes) |
De nombreux autres paramètres de configuration sont disponibles, mais assurez-vous de lire attentivement la documentation d'Apache Kafka avant de les expérimenter.
Partitionnement
Dans Kafka, les partitions constituent l'unité d'évolutivité. Par conséquent, le partitionnement est un moyen efficace d'augmenter votre débit, car il permet aux données thématiques de circuler en plusieurs flux parallèles.
Lorsqu'un producteur publie un message sur une rubrique, il peut choisir la partition à utiliser. Si l'ordre est important, notez qu'une partition est une séquence ordonnée d'enregistrements, mais qu'un thème comprend une ou plusieurs partitions. Si vous voulez qu'une série de messages soit distribuée dans l'ordre, assurez-vous que les messages se trouvent tous dans la même partition. Le moyen le plus simple pour obtenir ce résultat consiste à attribuer la même clé à tous ces messages.
Le producteur peut explicitement indiquer un numéro de partition lorsqu'il publie un message. Cela permet un contrôle direct, mais rend le code du producteur plus complexe car il prend la responsabilité de contrôler la sélection des partitions. Pour plus d'informations, voir l'appel de méthode Producer.partitionsFor. Par exemple, l'appel est décrit pour Kafka version 2.2.0.
Si le producteur n'indique pas de numéro de partition, la sélection de la partition s'effectue via un programme de partitionnement. Le programme de partitionnement par défaut intégré au producteur Kafka fonctionne comme suit :
- Si l'enregistrement n'a pas de clé, la sélection s'effectue de façon circulaire.
- Si l'enregistrement a une clé, la sélection s'effectue en calculant une valeur de hachage pour la clé. Ainsi, la même partition est sélectionnée pour tous les messages ayant la même clé.
Vous pouvez également écrire votre propre programme de partitionnement personnalisé. Un programme de partitionnement personnalisé peut choisir n'importe quel schéma pour affecter des enregistrements à des partitions. Par exemple, utiliser uniquement un sous-ensemble des informations de la clé ou un identificateur propre à l'application.
Ordre des messages
En général, Kafka écrit les messages dans l'ordre où le producteur les envoie. Toutefois, dans certaines situations, les tentatives peuvent entraîner la duplication ou la réorganisation des messages. Si vous souhaitez qu'une séquence de messages soit envoyée dans l'ordre, il est important de veiller à ce qu'ils soient tous écrits sur la même partition, car c'est le seul moyen de garantir l'ordre des messages.
Le producteur est également en mesure de réessayer d'envoyer des messages automatiquement. C'est une bonne idée d'activer cette fonction de relance, car l'alternative est que votre code d'application doit effectuer toutes les relances lui-même. La combinaison de la création de lots dans Kafka et des relances automatiques peut entraîner la duplication de messages et leur réorganisation.
Par exemple, si vous publiez une séquence de trois messages <M1, M2, M3> sur une rubrique. Les enregistrements peuvent tous faire partie d'un même lot, ils sont donc envoyés ensemble au responsable de la partition. Le responsable les écrit alors dans la partition et les réplique sous forme d'enregistrements distincts. En cas d'échec, il est possible que M1 et M2 soient ajoutés à la partition, mais pas M3. Le producteur ne reçoit pas d'accusé de réception, il réessaie donc d'envoyer <M1, M2, M3>. Le nouveau leader écrit M1, M2 et M3 sur la partition, qui contient maintenant <M1, M2, M1, M2, M3>, où le M1 dupliqué suit le M2 original. Si vous limitez à une le nombre de demandes en cours vers chaque courtier, vous évitez cette réorganisation. Il se peut qu'un enregistrement unique soit dupliqué, par exemple <M1, M2, M2, M3>, mais vous n'obtenez jamais de séquences hors séquence. Dans la version 0.11 ou ultérieure de Kafka, vous pouvez également utiliser la fonction de producteur idempotent pour éviter la duplication de M2
Avec Kafka, il est normal d'écrire les applications pour gérer les doublons occasionnels de messages, car l'impact sur les performances d'une seule requête en vol est significatif.
Accusés de réception des messages
Lorsque vous publiez un message, vous pouvez choisir le niveau d'accusés de réception requis en utilisant la configuration acks producer. Ce choix s'effectue dans un souci d'équilibre entre débit et fiabilité. Les trois niveaux
suivants existent.
- acks=0 (fiabilité réduite)
- Le message est considéré comme envoyé dès qu'il a été écrit sur le réseau. Aucun accusé de réception n'est attendu du responsable de la partition. Par conséquent, des messages peuvent se perdre si le responsable de la partition change. Ce niveau d'accusé de réception est rapide, mais il s'accompagne d'un risque de perte de message dans certaines situations.
- acks=1 (valeur par défaut)
- Le producteur accuse réception du message dès que le chef de partition a écrit avec succès son enregistrement dans la partition. Étant donné que l'accusé de réception intervient avant que l'enregistrement n'ait atteint les répliques synchronisées, le message peut être perdu si le leader tombe en panne, mais que les suiveurs n'ont pas encore reçu le message. Si la direction de la partition change, l'ancien chef en informe le producteur, qui peut traiter l'erreur et réessayer d'envoyer le message au nouveau chef. Comme les messages sont accusés avant que leur réception n'ait été confirmée par toutes les répliques, les messages qui ont été accusés mais qui n'ont pas encore été entièrement répliqués peuvent être perdus si la direction de la partition change.
- acks=all (fiabilité maximale)
- Le producteur accuse réception du message lorsque le chef de partition a réussi à écrire son enregistrement et que toutes les répliques synchronisées ont fait de même. Le message n'est pas perdu si le responsable de la partition change, sous réserve qu'au moins une réplique totalement synchronisée soit disponible.
Même si vous n'attendez pas que l'accusé de réception des messages soit envoyé au producteur, les messages sont toujours uniquement disponible pour consommation une fois validés, ce qui signifie que la réplication dans les répliques totalement synchronisées a été effectuée. En d'autres termes, le temps d'attente d'envoi des messages du point de vue du producteur est inférieur au temps d'attente de bout en bout mesuré du producteur qui envoie un message au consommateur qui reçoit le message.
Dans la mesure du possible, évitez d'attendre l'accusé de réception d'un message avant de publier le message suivant. Cette attente empêche le producteur de regrouper les messages en lots et réduit le débit de publication des messages jusqu'à atteindre le temps d'attente aller-retour du réseau.
Création de lots, régulation et compression
Par souci d'efficacité, le producteur rassemble des lots d'enregistrements pour les envoyer aux serveurs. Si vous activez la compression, le producteur compresse chaque lot, ce qui permet d'améliorer les performances puisque la quantité de données transférées sur le réseau est moindre.
Si vous tentez de publier des messages plus rapidement qu'ils peuvent être envoyés à un serveur, le producteur les place automatiquement en mémoire tampon dans des demandes par lots. Le producteur gère une mémoire tampon des enregistrements non envoyés pour chaque partition. Il arrive un moment où même la mise en lots ne permet pas d'atteindre le taux souhaité.
Un autre facteur a un impact important. Pour empêcher des producteurs ou consommateurs individuels d'envahir le cluster, Event Streams applique des quotas de débit. Le débit auquel chaque producteur envoie des données est calculé et tout producteur qui tente de dépasser son quota est régulé. La régulation est appliquée en reportant légèrement l'envoi des réponses au producteur. Ce processus fait généralement office de frein naturel.
Pour plus d'informations sur les conseils de débit, voir Limites et quotas.
En résumé, lorsqu'un message est publié, son enregistrement est d'abord écrit dans une mémoire tampon au niveau du producteur. En arrière-plan, le producteur crée des lots d'enregistrements et les envoie au serveur. Le serveur répond alors au producteur, en appliquant éventuellement un délai de régulation si le producteur publie trop vite. Si la mémoire tampon du producteur se remplit, l'appel d'envoi du producteur est retardé, mais peut finalement échouer avec une exception.
Niveaux de distribution
Kafka offre plusieurs niveaux de distribution de messages :
- Au maximum une fois: Les messages peuvent se perdre et ne sont pas redistribués.
- Au moins une fois: Les messages ne sont jamais perdus mais il peut y avoir des doublons.
- Exactement une fois: Les messages ne sont jamais perdus et il n'y a pas de doublons.
Les niveaux de distribution sont déterminés par les paramètres suivants :
acksretriesenable.idempotence
Par défaut, Kafka utilise au moins une fois la sémantique.
Pour activer la sémantique " exactement une fois", vous devez utiliser les producteurs idempotents ou transactionnels. Le producteur idempotent peut être activé en définissant enable.idempotence sur true et il garantit qu'une seule copie de chaque message est écrit à Kafka, même en cas de nouvel essai. Le producteur transactionnel, quant à lui, permet d'envoyer des données vers plusieurs partitions afin que tous les messages soient envoyés
avec succès, ou aucun d'entre eux. Autrement dit, une transaction est soit entièrement validée, soit entièrement annulée. Vous pouvez également inclure des décalages dans les transactions pour créer des applications qui lisent, traitent et
écrivent des messages dans Kafka.
Fragments de code
Ces fragments de code sont à un niveau élevé pour illustrer les concepts impliqués. Pour des exemples complets, voir les exemples Event Streams dans GitHub.
Pour connecter un consommateur à Event Streams, vous devez créer des données d'identification de service. Pour plus d'informations, voir Connexion à Event Streams.
Dans le code du producteur, vous devez d'abord générer l'ensemble des propriétés de configuration. Toutes les connexions à Event Streams sont sécurisées par TLS et l'authentification de l'utilisateur et du mot de passe, vous avez donc besoin de ces propriétés au minimum. Remplacez BOOTSTRAP_ENDPOINTS, USER et PASSWORD par les informations d'identification de votre propre service :
Properties props = new Properties();
props.put("bootstrap.servers", BOOTSTRAP_ENDPOINTS);
props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username=\"USER\" password=\"PASSWORD\";");
props.put("security.protocol", "SASL_SSL");
props.put("sasl.mechanism", "PLAIN");
props.put("ssl.protocol", "TLSv1.2");
props.put("ssl.enabled.protocols", "TLSv1.2");
props.put("ssl.endpoint.identification.algorithm", "HTTPS");
Pour envoyer des messages, vous devez également spécifier des sérialisateurs pour les clés et les valeurs, par exemple :
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
Ces sérialiseurs doivent correspondre aux désérialiseurs utilisés par les consommateurs.
Ensuite, utilisez un KafkaProducer pour envoyer des messages, où chaque message est représenté par un ProducerRecord. N'oubliez pas de fermer le KafkaProducer lorsque vous avez terminé. Ce code envoie seulement le message ; il n'attend pas pour
vérifier si l'envoi a réussi. Le message est envoyé à la rubrique T1, avec la clé key et la valeur value.
Producer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<String, String>("T1", "key", "value"));
producer.close();
La méthode send() est asynchrone et renvoie un Future que vous pouvez utiliser pour vérifier son achèvement :
Future<RecordMetadata> f = producer.send(new ProducerRecord<String, String>("T1", "key", "value"));
// Do some other stuff
// Now wait for the result of the send
RecordMetadata rm = f.get();
long offset = rm.offset;
Vous pouvez également fournir un rappel lors de l'envoi du message :
producer.send(new ProducerRecord<String,String>("T1","key","value", new Callback() {
public void onCompletion(RecordMetadata metadata, Exception exception) {
// This is called when the send completes, either successfully or with an exception
}
});
Pour plus d'informations, voir le Javadoc du client Kafka.{: external}.