Reprocessing des événements reçus de Kafka

Reprocessing des événements reçus de Kafka

Bonjour, Habr.

Récemment, j'ai partagé mon expérience sur les paramÚtres que nous utilisons le plus souvent dans notre équipe pour Kafka Producer et Consumer, afin de nous rapprocher d'une livraison garantie. Dans cet article, je souhaite expliquer comment nous avons organisé le traitement des événements reçus de Kafka en raison de l'indisponibilité temporaire d'un systÚme externe.

Les applications modernes fonctionnent dans un environnement trÚs complexe. La logique métier, enveloppée dans une pile technologique moderne, fonctionne dans une image Docker gérée par un orchestrateur tel que Kubernetes ou OpenShift, et communique avec d'autres applications ou des solutions d'entreprise via une chaßne de routeurs physiques et virtuels. Dans un tel environnement, quelque chose peut toujours casser, c'est pourquoi le traitement des événements en cas d'indisponibilité de l'un des systÚmes externes est une partie importante de nos processus métier.

Comment c'était avant Kafka

Auparavant, dans le projet, nous utilisions IBM MQ pour la livraison asynchrone des messages. En cas d'erreur survenant dans le processus de travail du service, le message reçu pouvait ĂȘtre placĂ© dans une dead-letter-queue (DLQ) pour un examen manuel ultĂ©rieur. La DLQ Ă©tait créée Ă  cĂŽtĂ© de la file d'attente entrante, le transfert du message se faisait au sein d'IBM MQ.

Si l'erreur Ă©tait temporaire et que nous pouvions le dĂ©terminer (par exemple, ResourceAccessException lors d'un appel HTTP ou MongoTimeoutException lors d'une requĂȘte Ă  MongoDb), alors la stratĂ©gie des appels rĂ©pĂ©tĂ©s entrait en vigueur. Peu importe la branche de la logique de l'application, le message d'origine Ă©tait transfĂ©rĂ© soit dans une file d'attente systĂšme pour un envoi diffĂ©rĂ©, soit dans une application distincte qui avait Ă©tĂ© créée il y a longtemps pour renvoyer des messages. Dans ce cas, le numĂ©ro de rĂ©expĂ©dition, liĂ© Ă  l'intervalle de retard ou Ă  la fin de la stratĂ©gie au niveau de l'application, Ă©tait enregistrĂ© dans l'en-tĂȘte du message. Si nous atteignons la fin de la stratĂ©gie mais que le systĂšme externe reste indisponible, alors le message sera placĂ© dans une DLQ pour un examen manuel.

Recherche de solution

En cherchant sur Internet, on peut trouver ce qui suit la solution. En bref, il est proposĂ© de crĂ©er un topic par intervalle de retard et de mettre en Ɠuvre des Consumers du cĂŽtĂ© de l'application, qui liront les messages avec le retard nĂ©cessaire.

Reprocessing des événements reçus de Kafka

Bien qu'il y ait de nombreux avis positifs, il me semble qu'il n'est pas tout Ă  fait rĂ©ussi. Tout d'abord, parce que le dĂ©veloppeur, en plus de rĂ©pondre aux exigences commerciales, devra consacrer beaucoup de temps Ă  la mise en Ɠuvre du mĂ©canisme dĂ©crit.

De plus, si la gestion des accĂšs est activĂ©e sur le cluster Kafka, il faudra passer un certain temps Ă  crĂ©er des topics et Ă  garantir les accĂšs nĂ©cessaires Ă  ceux-ci. En outre, il sera nĂ©cessaire de choisir le bon paramĂštre retention.ms pour chacun des topics de rĂ©essai, afin que les messages puissent ĂȘtre renvoyĂ©s Ă  temps et ne disparaissent pas. La mise en Ɠuvre et la demande d'accĂšs devront ĂȘtre rĂ©pĂ©tĂ©es pour chaque service existant ou nouveau.

Voyons maintenant quels mĂ©canismes de traitement des messages en double nous propose spring en gĂ©nĂ©ral et spring-kafka en particulier. Spring-kafka a une dĂ©pendance transitive sur spring-retry, qui fournit des abstractions pour gĂ©rer diffĂ©rentes BackOffPolicy. C'est un outil assez flexible, mais son inconvĂ©nient majeur est le stockage des messages Ă  renvoyer en mĂ©moire de l'application. Cela signifie qu'un redĂ©marrage de l'application en raison d'une mise Ă  jour ou d'une erreur pendant l'exploitation entraĂźnera la perte de tous les messages en attente de traitement Ă  nouveau. Étant donnĂ© que ce point est critique pour notre systĂšme, nous avons choisi de ne pas le considĂ©rer davantage.

La bibliothÚque spring-kafka propose plusieurs implémentations de ContainerAwareErrorHandler, par exemple SeekToCurrentErrorHandler, qui permet de traiter le message plus tard sans déplacer l'offset en cas d'erreur. Depuis la version 2.3 de spring-kafka, il est possible de définir une BackOffPolicy.

Cette approche permet aux messages en réessai de survivre au redémarrage de l'application, mais le mécanisme DLQ est toujours absent. C'est cette option que nous avons choisie au début de 2019, en pensant de maniÚre optimiste que le DLQ ne serait pas nécessaire (nous avons eu de la chance et il ne fut effectivement pas nécessaire pendant quelques mois d'exploitation de l'application avec ce systÚme de réessai). Des erreurs temporaires provoquaient l'activation de SeekToCurrentErrorHandler. Les autres erreurs étaient enregistrées dans les logs, entraßnant un décalage de l'offset, et le traitement se poursuivait avec le message suivant.

Solution finale

La mise en Ɠuvre basĂ©e sur SeekToCurrentErrorHandler nous a poussĂ©s Ă  dĂ©velopper notre propre mĂ©canisme pour le renvoi des messages.

Avant tout, nous voulions tirer parti de l'expĂ©rience existante et l'Ă©largir en fonction de la logique de l'application. Pour une application Ă  logique linĂ©aire, il serait optimal d'arrĂȘter la lecture de nouveaux messages pendant un court laps de temps dĂ©fini dans la stratĂ©gie de rĂ©essai. Pour les autres applications, nous souhaiterions avoir un point central qui garantirait l'application de la stratĂ©gie de rĂ©essai. De plus, ce point unique doit disposer d'une fonctionnalitĂ© DLQ pour les deux approches.

La stratĂ©gie de rĂ©essai elle-mĂȘme doit ĂȘtre stockĂ©e dans l'application responsable de la rĂ©cupĂ©ration du prochain intervalle en cas d'erreur temporaire.

ArrĂȘt du Consumer pour une application Ă  logique linĂ©aire

Lors de l'utilisation de spring-kafka, le code pour arrĂȘter le Consumer peut ressembler Ă  ceci :

public void pauseListenerContainer(MessageListenerContainer listenerContainer, 
                                   Instant retryAt) {
        if (nonNull(retryAt) && listenerContainer.isRunning()) {
            listenerContainer.stop();
            taskScheduler.schedule(() -> listenerContainer.start(), retryAt);
            return;
        }
        // vers DLQ
    }

Dans cet exemple, retryAt est le moment auquel il faut redémarrer le MessageListenerContainer s'il est encore en cours d'exécution. Le redémarrage se fera dans un thread séparé, lancé par TaskScheduler, dont l'implémentation est également fournie par spring.

Nous trouvons la valeur de retryAt de la maniĂšre suivante :

  1. Nous recherchons la valeur du compteur de réessais.
  2. En fonction de la valeur du compteur, l'intervalle de dĂ©lai actuel est recherchĂ© dans la stratĂ©gie de rĂ©essai. La stratĂ©gie est dĂ©clarĂ©e dans l'application elle-mĂȘme, et nous avons choisi le format JSON pour son stockage.
  3. L'intervalle trouvé dans le tableau JSON contient le nombre de secondes aprÚs lesquelles il faudra réessayer le traitement. Ce nombre de secondes est ajouté au temps actuel, formant ainsi la valeur pour retryAt.
  4. Si l'intervalle n'est pas trouvé, la valeur de retryAt est nulle et le message sera envoyé dans la DLQ pour un traitement manuel.

Avec cette approche, il ne reste qu'à conserver le nombre de tentatives pour chaque message actuellement en cours de traitement, par exemple dans la mémoire de l'application. Conserver le compteur de tentatives en mémoire n'est pas critique pour cette approche, car une application avec une logique linéaire ne peut pas traiter l'ensemble. Contrairement à spring-retry, le redémarrage de l'application ne conduira pas à la perte de tous les messages à réitérer, mais simplement à un redémarrage de la stratégie.

Cette approche permet de rĂ©duire la charge sur le systĂšme externe, qui peut ĂȘtre indisponible en raison d'une trĂšs forte charge. En d'autres termes, en plus de la rĂ©itĂ©ration, nous avons rĂ©ussi Ă  mettre en Ɠuvre le modĂšle. disjoncteur.

Dans notre cas, le seuil d'erreur est de seulement 1, et pour minimiser le temps d'arrĂȘt du systĂšme en raison d'une panne de rĂ©seau temporaire, nous utilisons une stratĂ©gie de rĂ©itĂ©ration trĂšs granulaire avec de courtes intervalles de dĂ©lai. Cela peut ne pas convenir Ă  toutes les applications du groupe, donc le ratio entre le seuil d'erreur et la taille de l'intervalle doit ĂȘtre ajustĂ© en fonction des caractĂ©ristiques du systĂšme.

Une application distincte pour le traitement des messages provenant d'applications avec une logique non déterministe

Voici un exemple de code qui envoie un message Ă  une telle application (Retryer), qui effectuera une nouvelle tentative d'envoi vers le sujet DESTINATION lorsque le temps RETRY_AT sera atteint :


public  void retry(ConsumerRecord record, String retryToTopic, 
                         Instant retryAt, String counter, String groupId, Exception e) {
        Headers headers = ofNullable(record.headers()).orElse(new RecordHeaders());
        List
arrayOfHeaders = new ArrayList(Arrays.asList(headers.toArray())); updateHeader(arrayOfHeaders, GROUP_ID, groupId::getBytes); updateHeader(arrayOfHeaders, DESTINATION, retryToTopic::getBytes); updateHeader(arrayOfHeaders, ORIGINAL_PARTITION, () -> Integer.toString(record.partition()).getBytes()); if (nonNull(retryAt)) { updateHeader(arrayOfHeaders, COUNTER, counter::getBytes); updateHeader(arrayOfHeaders, SEND_TO, "retry"::getBytes); updateHeader(arrayOfHeaders, RETRY_AT, retryAt.toString()::getBytes); } else { updateHeader(arrayOfHeaders, REASON, ExceptionUtils.getStackTrace(e)::getBytes); updateHeader(arrayOfHeaders, SEND_TO, "backout"::getBytes); } ProducerRecord messageToSend = new ProducerRecord(retryTopic, null, null, record.key(), record.value(), arrayOfHeaders); kafkaTemplate.send(messageToSend); }

L'exemple montre qu'une grande quantitĂ© d'informations est transmise dans les en-tĂȘtes. La valeur RETRY_AT est dĂ©terminĂ©e de la mĂȘme maniĂšre que pour le mĂ©canisme de rĂ©pĂ©tition via l'arrĂȘt du Consumer. En plus de DESTINATION et RETRY_AT, nous transmettons :

  • GROUP_ID, par lequel nous groupons les messages pour une analyse manuelle et simplifions la recherche.
  • ORIGINAL_PARTITION, afin d'essayer de conserver le mĂȘme Consumer pour le retraitement. Ce paramĂštre peut ĂȘtre null, auquel cas une nouvelle partition sera obtenue par la clĂ© record.key() du message original.
  • La valeur mise Ă  jour du COUNTER, pour suivre la stratĂ©gie de rĂ©pĂ©titions.
  • SEND_TO — une constante qui indique s'il faut renvoyer le message pour retraitement lorsqu'on atteint RETRY_AT ou le placer dans la DLQ.
  • REASON — la raison pour laquelle le traitement du message a Ă©tĂ© interrompu.

Le Retryer conserve les messages pour un nouvel envoi et une analyse manuelle dans PostgreSQL. Une tùche est lancée selon un minuteur, qui trouve les messages dont le RETRY_AT est atteint et les renvoie dans la partition ORIGINAL_PARTITION du sujet DESTINATION avec la clé record.key().

AprÚs l'envoi, les messages sont supprimés de PostgreSQL. L'analyse manuelle des messages se déroule dans une interface simple qui interagit avec le Retryer via l'API REST. Ses principales caractéristiques sont la réexpédition ou la suppression de messages de la DLQ, la consultation des informations d'erreur et la recherche de messages, par exemple, par nom d'erreur.

Comme l'accÚs est contrÎlé sur nos clusters, il est nécessaire de demander l'accÚs au sujet écouté par le Retryer et de permettre au Retryer d'écrire dans le sujet DESTINATION. C'est peu pratique, mais, contrairement à l'approche avec un sujet basé sur des intervalles, nous disposons d'une véritable DLQ et d'une interface utilisateur pour la gérer.

Il arrive que le sujet entrant soit lu par plusieurs groupes de consommateurs diffĂ©rents, dont les applications mettent en Ɠuvre une logique variĂ©e. La rĂ©pĂ©tition d'un message via le Retryer pour l'une de ces applications entraĂźnera un doublon pour l'autre. Pour se protĂ©ger contre cela, nous crĂ©ons un sujet distinct pour le retraitement. Le sujet entrant et le sujet de rĂ©pĂ©tition peuvent ĂȘtre lus par le mĂȘme Consumer sans aucune restriction.

Reprocessing des événements reçus de Kafka

Par dĂ©faut, cette approche ne fournit pas la possibilitĂ© d'un circuit breaker, mais cela peut ĂȘtre ajoutĂ© dans l'application Ă  l'aide de spring-cloud-netflix ou du nouveau spring cloud circuit breaker, en enveloppant les lieux d'appels aux services externes dans les abstractions appropriĂ©es. De plus, cela permet de choisir une stratĂ©gie pour bulkhead un modĂšle, ce qui peut Ă©galement ĂȘtre utile. Par exemple, dans spring-cloud-netflix, cela peut ĂȘtre une pool de threads ou un sĂ©maphore.

Sortie

En conséquence, nous avons développé une application distincte qui permet de répéter le traitement d'un message en cas d'indisponibilité temporaire d'un systÚme externe.

L'un des principaux avantages de l'application est qu'elle peut ĂȘtre utilisĂ©e par des systĂšmes externes opĂ©rant sur le mĂȘme cluster Kafka, sans nĂ©cessiter d'importantes modifications de leur part ! Cette application n'aura besoin que d'accĂ©der au topic retry, de remplir quelques en-tĂȘtes Kafka et d'envoyer le message au Retryer. Il n'est pas nĂ©cessaire de dĂ©ployer d'infrastructure supplĂ©mentaire. De plus, pour rĂ©duire le nombre de messages transfĂ©rĂ©s entre l'application et le Retryer, nous avons isolĂ© des applications avec une logique linĂ©aire et mis en place un traitement rĂ©pĂ©tĂ© via l'arrĂȘt du Consumer.

Source : habr.com

Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS đŸ”„ Acheter un hĂ©bergement fiable pour les sites avec protection DDoS, serveurs VPS VDS | ProHoster