{"id":56186,"date":"2020-02-06T00:00:00","date_gmt":"2020-02-05T21:00:00","guid":{"rendered":"https:\/\/prohoster.info\/blog\/blog_prohoster\/povtornaya-obrabotka-sobytij-poluchennyh-iz-kafka"},"modified":"2020-02-18T14:04:24","modified_gmt":"2020-02-18T11:04:24","slug":"povtornaya-obrabotka-sobytij-poluchennyh-iz-kafka","status":"publish","type":"post","link":"https:\/\/prohoster.info\/ro\/blog\/administrirovanie\/povtornaya-obrabotka-sobytij-poluchennyh-iz-kafka","title":{"rendered":"Reprelucrarea evenimentelor primite din Kafka","gt_translate_keys":[{"key":"rendered","format":"text"}]},"content":{"rendered":"<p><img decoding=\"async\" alt=\"Reprelucrarea evenimentelor primite din Kafka\" src=\"\/wp-content\/uploads\/2020\/02\/d04319691af5c5de40ccfe4a2268558c.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<p><\/p>\n<p>Bun\u0103, Habr.<\/p>\n<p><\/p>\n<p>Recent am <noindex><a rel=\"nofollow\" href=\"https:\/\/habr.com\/ru\/company\/tinkoff\/blog\/481784\/\">\u00eemp\u0103rt\u0103\u0219it experien\u021ba<\/a><\/noindex> despre ce parametrii folosim cel mai adesea \u00een echip\u0103 pentru Kafka Producer \u0219i Consumer pentru a ne apropia de livrarea garantat\u0103. \u00cen acest articol vreau s\u0103 discut despre cum am organizat procesarea repetat\u0103 a evenimentului primit din Kafka, ca urmare a indisponibilit\u0103\u021bii temporare a unei sisteme externe.<\/p>\n<p><\/p>\n<p>Aplica\u021biile moderne func\u021bioneaz\u0103 \u00eentr-un mediu foarte complex. Logica de afaceri, \u00eenv\u0103luit\u0103 \u00een o stiv\u0103 tehnologic\u0103 modern\u0103, func\u021bioneaz\u0103 \u00eentr-o imagine Docker, gestionat\u0103 de un orchestrator precum Kubernetes sau OpenShift, \u0219i comunic\u0103 cu alte aplica\u021bii sau solu\u021bii enterprise printr-o serie de routere fizice \u0219i virtuale. \u00centr-un astfel de mediu, mereu poate ap\u0103rea ceva care se stric\u0103, astfel c\u0103 procesarea repetat\u0103 a evenimentelor \u00een cazul \u00een care una dintre sistemele externe devine indisponibil\u0103 este o parte important\u0103 a proceselor noastre de afaceri.<\/p>\n<p><noindex><a rel=\"nofollow\" name=\"habracut\"><\/a><\/noindex><\/p>\n<h2 id=\"kak-bylo-do-kafka\">Cum era \u00eenainte de Kafka<\/h2>\n<p><\/p>\n<p>Anterior, \u00een proiect, foloseam IBM MQ pentru livrarea asincron\u0103 a mesajelor. C\u00e2nd ap\u0103rea vreo eroare \u00een procesul de func\u021bionare a serviciului, mesajul primit putea fi plasat \u00eentr-o coad\u0103 de mesaje moarte (DLQ) pentru o analiz\u0103 ulterioar\u0103 manual\u0103. DLQ era creat\u0103 l\u00e2ng\u0103 coada de intrare, mutarea mesajului av\u00e2nd loc \u00een interiorul IBM MQ. <\/p>\n<p><\/p>\n<p>Dac\u0103 eroarea avea un caracter temporar \u0219i puteam determina asta (de exemplu, ResourceAccessException \u00een timpul unui apel HTTP sau MongoTimeoutException \u00een timpul unei interog\u0103ri \u00een MongoDb), atunci intervenea strategia de apeluri repetate. Indiferent de ramificarea logicii aplica\u021biei, mesajul original era mutat fie \u00een coada sistemului pentru trimitere \u00eent\u00e2rziat\u0103, fie \u00een o aplica\u021bie separat\u0103, care a fost creat\u0103 cu mult timp \u00een urm\u0103 pentru retrimiterea mesajelor. \u00cen acest context, \u00een antetul mesajului se \u00eenregistra num\u0103rul de retrimitere, care era legat de intervalul de \u00eent\u00e2rziere sau de sf\u00e2r\u0219itul strategiei la nivel de aplica\u021bie. Dac\u0103 am atins sf\u00e2r\u0219itul strategiei, dar sistemul extern era \u00eenc\u0103 indisponibil, mesajul era plasat \u00een DLQ pentru analiz\u0103 manual\u0103.<\/p>\n<p><\/p>\n<h2 id=\"poisk-resheniya\">C\u0103utarea solu\u021biei<\/h2>\n<p><\/p>\n<p><noindex><a rel=\"nofollow\" href=\"https:\/\/www.google.com\/search?q=kafka+retry+message\">C\u0103ut\u00e2nd pe internet<\/a><\/noindex>, po\u021bi g\u0103si urm\u0103toarele <noindex><a rel=\"nofollow\" href=\"https:\/\/blog.pragmatists.com\/retrying-consumer-architecture-in-the-apache-kafka-939ac4cb851a\">solu\u021bie<\/a><\/noindex>. Pe scurt, se propune s\u0103 se creeze o tem\u0103 pentru fiecare interval de \u00eent\u00e2rziere \u0219i s\u0103 se implementeze de partea aplica\u021biei Consumers care vor citi mesajele cu \u00eent\u00e2rzierea necesar\u0103. <\/p>\n<p><\/p>\n<p><img decoding=\"async\" alt=\"Reprelucrarea evenimentelor primite din Kafka\" src=\"\/wp-content\/uploads\/2020\/02\/718c4c3f71ccf95b5a6f60e1b2c685d2.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<p><\/p>\n<p>\u00cen ciuda numeroaselor recenzii pozitive, mi se pare c\u0103 nu este complet reu\u0219it. \u00cen primul r\u00e2nd, deoarece dezvoltatorul, pe l\u00e2ng\u0103 implementarea cerin\u021belor de afaceri, va trebui s\u0103 aloce mult timp pentru realizarea mecanismului descris.<\/p>\n<p><\/p>\n<p>\u00cen plus, dac\u0103 pe clusterul Kafka este activat\u0103 gestionarea accesului, va trebui s\u0103 aloce ceva timp pentru a crea topicuri \u0219i a asigura accesul necesar la acestea. \u00cen plus, va trebui s\u0103 selecteze parametrul corect retention.ms pentru fiecare din topicurile de retry, astfel \u00eenc\u00e2t mesajele s\u0103 fie retransmise \u00een timp \u0219i s\u0103 nu dispar\u0103 din acestea. Implementarea \u0219i solicitarea accesului va trebui repetat\u0103 pentru fiecare serviciu existent sau nou.<\/p>\n<p><\/p>\n<p>S\u0103 vedem acum ce mecanisme de procesare repetat\u0103 a mesajelor ne ofer\u0103 spring \u00een general \u0219i spring-kafka \u00een special. Spring-kafka are o dependen\u021b\u0103 tranzitiv\u0103 de spring-retry, care ofer\u0103 abstra\u021bii pentru gestionarea diferitelor BackOffPolicy. Acesta este un instrument destul de flexibil, dar dezavantajul s\u0103u semnificativ este stocarea mesajelor pentru retransmitere \u00een memoria aplica\u021biei. Asta \u00eenseamn\u0103 c\u0103 repornirea aplica\u021biei din cauza unei actualiz\u0103ri sau a unei erori \u00een timpul utiliz\u0103rii va duce la pierderea tuturor mesajelor care a\u0219teapt\u0103 procesarea repetat\u0103. Deoarece acest punct este critic pentru sistemul nostru, nu l-am mai analizat ulterior.<\/p>\n<p><\/p>\n<p>Spring-kafka \u00een sine ofer\u0103 mai multe implement\u0103ri ale ContainerAwareErrorHandler, de exemplu, <noindex><a rel=\"nofollow\" href=\"https:\/\/github.com\/spring-projects\/spring-kafka\/blob\/master\/spring-kafka\/src\/main\/java\/org\/springframework\/kafka\/listener\/SeekToCurrentErrorHandler.java\">SeekToCurrentErrorHandler<\/a><\/noindex>, cu ajutorul c\u0103ruia se poate, f\u0103r\u0103 a muta offset-ul \u00een caz de eroare, procesa mesajul mai t\u00e2rziu. \u00cencep\u00e2nd cu versiunea spring-kafka 2.3, a ap\u0103rut posibilitatea de a defini BackOffPolicy.<\/p>\n<p><\/p>\n<p>Aceast\u0103 abordare permite mesajelor retrase s\u0103 supravie\u021buiasc\u0103 repornirii aplica\u021biei, dar mecanismul DLQ este \u00eenc\u0103 absent. Acesta a fost variant\u0103 pe care am ales-o la \u00eenceputul anului 2019, fiind optim\u0219ti c\u0103 DLQ nu va fi necesar (am avut noroc \u0219i \u00eentr-adev\u0103r nu a fost nevoie \u00een c\u00e2teva luni de utilizare a aplica\u021biei cu un astfel de sistem de procesare repetat\u0103). Erorile temporare au dus la activarea SeekToCurrentErrorHandler. Celelalte erori au fost tip\u0103rite \u00een loguri, au dus la mutarea offset-ului \u0219i procesarea a continuat cu urm\u0103torul mesaj.<\/p>\n<p><\/p>\n<h2 id=\"itogovoe-reshenie\">Solu\u021bia final\u0103<\/h2>\n<p><\/p>\n<p>Implementarea bazat\u0103 pe SeekToCurrentErrorHandler ne-a \u00eempins s\u0103 dezvolt\u0103m un mecanism propriu pentru retransmiterea mesajelor.<\/p>\n<p><\/p>\n<p>\u00cen primul r\u00e2nd, am dorit s\u0103 folosim experien\u021ba existent\u0103 \u0219i s\u0103 o extindem \u00een func\u021bie de logica aplica\u021biei. Pentru o aplica\u021bie cu o logic\u0103 liniar\u0103, ar fi optim s\u0103 oprim citirea de noi mesaje pe o perioad\u0103 scurt\u0103 de timp, specificat\u0103 \u00een cadrul strategiei de retransmitere. Pentru celelalte aplica\u021bii, am dorit s\u0103 avem un punct unic care s\u0103 asigure implementarea strategiei de retransmitere. \u00cen plus, acest punct unic ar trebui s\u0103 dispun\u0103 de func\u021bionalitatea DLQ pentru ambele abord\u0103ri.<\/p>\n<p><\/p>\n<p>Strategia de retransmitere propriu-zis\u0103 ar trebui s\u0103 fie stocat\u0103 \u00een aplica\u021bia responsabil\u0103 pentru ob\u021binerea urm\u0103torului interval \u00een cazul unei erori temporare.<\/p>\n<p><\/p>\n<h3 id=\"ostanovka-consumera-dlya-prilozheniya-s-lineynoy-logikoy\">Oprirea Consumer-ului pentru o aplica\u021bie cu logic\u0103 liniar\u0103<\/h3>\n<p><\/p>\n<p>Atunci c\u00e2nd lucr\u0103m cu spring-kafka, codul pentru oprirea Consumer-ului ar putea ar\u0103ta cam a\u0219a:<\/p>\n<p><\/p>\n<pre><code class=\"java\">public void pauseListenerContainer(MessageListenerContainer listenerContainer, \n                                   Instant retryAt) {\n        if (nonNull(retryAt) &amp;&amp; listenerContainer.isRunning()) {\n            listenerContainer.stop();\n            taskScheduler.schedule(() -&gt; listenerContainer.start(), retryAt);\n            return;\n        }\n        \/\/ to DLQ\n    }<\/code><\/pre>\n<p><\/p>\n<p>\u00cen exemplu, retryAt este momentul \u00een care trebuie s\u0103 restart\u0103m MessageListenerContainer, dac\u0103 acesta este \u00eenc\u0103 activ. Restartarea va avea loc \u00eentr-un fir de execu\u021bie separat, lansat \u00een TaskScheduler, a c\u0103rei implementare o ofer\u0103 de asemenea spring. <\/p>\n<p><\/p>\n<p>Valoarea retryAt o g\u0103sim \u00een urm\u0103torul mod:<\/p>\n<p><\/p>\n<ol>\n<li>Se caut\u0103 valoarea contorului de retransmiteri.<\/li>\n<li>Conform valorii contorului, se caut\u0103 intervalul curent de \u00eent\u00e2rziere \u00een strategia de retransmitere. Strategia este declarat\u0103 \u00een cadrul aplica\u021biei, iar pentru stocarea acesteia am ales formatul JSON.<\/li>\n<li>Intervalul g\u0103sit \u00een array-ul JSON con\u021bine num\u0103rul de secunde dup\u0103 care trebuie s\u0103 repet\u0103m procesarea. Acest num\u0103r de secunde este ad\u0103ugat la timpul curent, form\u00e2nd valoarea pentru retryAt.<\/li>\n<li>Dac\u0103 intervalul nu este g\u0103sit, atunci valoarea retryAt este null \u0219i mesajul va fi trimis \u00een DLQ pentru o analiz\u0103 manual\u0103.<\/li>\n<\/ol>\n<p><\/p>\n<p>Cu aceast\u0103 abordare, trebuie doar s\u0103 p\u0103str\u0103m num\u0103rul apelurilor repetate pentru fiecare mesaj care este acum \u00een procesare, de exemplu, \u00een memoria aplica\u021biei. P\u0103strarea unui contor de \u00eencerc\u0103ri \u00een memorie nu este critic\u0103 pentru aceast\u0103 abordare, deoarece aplica\u021bia cu o logic\u0103 liniar\u0103 nu poate efectua procesarea \u00een \u00eentregime. Spre deosebire de spring-retry, repornirea aplica\u021biei nu va duce la pierderea tuturor mesajelor pentru re-procesare, ci pur \u0219i simplu la reluarea strategiei. <\/p>\n<p><\/p>\n<p>Aceast\u0103 abordare ajut\u0103 la reducerea sarcinii asupra sistemului extern, care poate fi inaccesibil din cauza unei sarcini foarte mari. Cu alte cuvinte, \u00een plus fa\u021b\u0103 de re-procesare, am realizat implementarea modelului <noindex><a rel=\"nofollow\" href=\"https:\/\/microservices.io\/patterns\/reliability\/circuit-breaker.html\">circuit breaker<\/a><\/noindex>.<\/p>\n<p><\/p>\n<p>\u00cen cazul nostru, pragul de eroare este doar 1, iar pentru a minimiza timpul de nefunc\u021bionare a sistemului din cauza unei \u00eentreruperi temporare a re\u021belei, folosim o strategie de apeluri repetitive foarte granular\u0103 cu intervale de \u00eent\u00e2rzieri mici. Acest lucru poate s\u0103 nu fie potrivit pentru toate aplica\u021biile grupului de companii, a\u0219a c\u0103 raportul \u00eentre pragul de eroare \u0219i dimensiunea intervalului trebuie ajustat, baz\u00e2ndu-se pe caracteristicile sistemului.<\/p>\n<p><\/p>\n<h3 id=\"otdelnoe-prilozhenie-dlya-obrabotki-soobscheniy-ot-prilozheniy-s-nedeterminirovannoy-logikoy\">O aplica\u021bie separat\u0103 pentru procesarea mesajelor de la aplica\u021bii cu o logic\u0103 nedeterminist\u0103<\/h3>\n<p><\/p>\n<p>Iat\u0103 un exemplu de cod care trimite un mesaj c\u0103tre o astfel de aplica\u021bie (Retryer), care va re\u00eencerca trimiterea \u00een topicul DESTINATION atunci c\u00e2nd se atinge timpul RETRY_AT:<\/p>\n<p><\/p>\n<pre><code class=\"java\">\npublic  void retry(ConsumerRecord record, String retryToTopic, \n                         Instant retryAt, String counter, String groupId, Exception e) {\n        Headers headers = ofNullable(record.headers()).orElse(new RecordHeaders());\n        List<Header> arrayOfHeaders = \n            new ArrayList(Arrays.asList(headers.toArray()));\n        updateHeader(arrayOfHeaders, GROUP_ID, groupId::getBytes);\n        updateHeader(arrayOfHeaders, DESTINATION, retryToTopic::getBytes);\n        updateHeader(arrayOfHeaders, ORIGINAL_PARTITION, \n                     () -&gt; Integer.toString(record.partition()).getBytes());\n        if (nonNull(retryAt)) {\n            updateHeader(arrayOfHeaders, COUNTER, counter::getBytes);\n            updateHeader(arrayOfHeaders, SEND_TO, \"retry\"::getBytes);\n            updateHeader(arrayOfHeaders, RETRY_AT, retryAt.toString()::getBytes);\n        } else {\n            updateHeader(arrayOfHeaders, REASON, \n                         ExceptionUtils.getStackTrace(e)::getBytes);\n            updateHeader(arrayOfHeaders, SEND_TO, \"backout\"::getBytes);\n        }\n        ProducerRecord messageToSend =\n            new ProducerRecord(retryTopic, null, null, record.key(), record.value(), arrayOfHeaders);\n        kafkaTemplate.send(messageToSend);\n    }<\/code><\/pre>\n<p><\/p>\n<p>Din exemplu, se observ\u0103 c\u0103 se transmite multe informa\u021bii \u00een header-e. Valoarea RETRY_AT se afl\u0103 la fel ca \u0219i pentru mecanismul de re\u00eencercare prin oprirea Consumer-ului. Pe l\u00e2ng\u0103 DESTINATION \u0219i RETRY_AT, mai transmitem:<\/p>\n<p><\/p>\n<ul>\n<li>GROUP_ID, care este utilizat pentru a grupa mesajele pentru analiza manual\u0103 \u0219i simplificarea c\u0103ut\u0103rii.<\/li>\n<li>ORIGINAL_PARTITION, pentru a \u00eencerca s\u0103 p\u0103str\u0103m acela\u0219i Consumer pentru procesarea ulterioar\u0103. Acest parametru poate fi egal cu null, caz \u00een care o nou\u0103 partition va fi ob\u021binut\u0103 pe baza cheii record.key() a mesajului original.<\/li>\n<li>Valoarea actualizat\u0103 COUNTER, pentru a respecta strategia de re\u00eencerc\u0103ri.<\/li>\n<li>SEND_TO - constant\u0103 care arat\u0103 dac\u0103 mesajul trebuie trimis pentru reprocesare la atingerea RETRY_AT sau s\u0103 fie plasat \u00een DLQ.<\/li>\n<li>REASON - motivul pentru care procesarea mesajului a fost oprit\u0103.<\/li>\n<\/ul>\n<p><\/p>\n<p>Retryer salveaz\u0103 mesajele pentru trimiterea ulterioar\u0103 \u0219i analiza manual\u0103 \u00een PostgreSQL. La un interval stabilit, o sarcin\u0103 este activat\u0103 care g\u0103se\u0219te mesajele cu RETRY_AT \u00eemplinit \u0219i le trimite \u00eenapoi \u00een partition-ul ORIGINAL_PARTITION al topic-ului DESTINATION cu cheia record.key().<\/p>\n<p><\/p>\n<p>Dup\u0103 trimiterea mesajului, acesta este \u0219ters din PostgreSQL. Analiza manual\u0103 a mesajelor se desf\u0103\u0219oar\u0103 \u00eentr-o interfa\u021b\u0103 simpl\u0103, care comunic\u0103 cu Retryer prin REST API. Principalele func\u021bii ale acesteia includ retransmiterea sau \u0219tergerea mesajelor din DLQ, vizualizarea informa\u021biilor despre erori \u0219i c\u0103utarea mesajelor, de exemplu, dup\u0103 numele erorii. <\/p>\n<p><\/p>\n<p>Dat fiind c\u0103 gestionarea accesului este activat\u0103 pe clusterele noastre, este necesar s\u0103 se solicite suplimentar drepturi de acces pentru topic-ul pe care \u00eel ascult\u0103 Retryer \u0219i s\u0103 se permit\u0103 Retryer-ului s\u0103 scrie \u00een topic-ul DESTINATION. Acest lucru este inconfortabil, dar, spre deosebire de abordarea cu topic-ul pe interval, ob\u021binem o DLQ complet func\u021bional\u0103 \u0219i o interfa\u021b\u0103 pentru gestionarea acesteia.<\/p>\n<p><\/p>\n<p>Exist\u0103 cazuri \u00een care topic-ul de intrare este citit de mai multe grupuri de consumatori diferite, ale c\u0103ror aplica\u021bii implementeaz\u0103 o logic\u0103 diferit\u0103. Reprocesarea unui mesaj prin Retryer pentru una dintre aceste aplica\u021bii va duce la un duplicat pe cealalt\u0103. Pentru a ne proteja de aceast\u0103 problem\u0103, cre\u0103m un topic separat pentru reprocesare. Topic-ul de intrare \u0219i topic-ul de retry pot fi citite de acela\u0219i Consumer f\u0103r\u0103 restric\u021bii. <\/p>\n<p><\/p>\n<p><img decoding=\"async\" alt=\"Reprelucrarea evenimentelor primite din Kafka\" src=\"\/wp-content\/uploads\/2020\/02\/06f2355839c6dd812ee3304bf2c72226.jpg\" style=\"display:block;margin: 0 auto;\" \/><\/p>\n<p><\/p>\n<p>\u00cen mod implicit, aceast\u0103 abordare nu ofer\u0103 posibilitatea unui circuit breaker, dar acesta poate fi ad\u0103ugat \u00een aplica\u021bie cu ajutorul <noindex><a rel=\"nofollow\" href=\"https:\/\/spring.io\/projects\/spring-cloud-netflix\">spring-cloud-netflix<\/a><\/noindex> sau noul <noindex><a rel=\"nofollow\" href=\"https:\/\/spring.io\/projects\/spring-cloud-circuitbreaker\">spring cloud circuit breaker<\/a><\/noindex>, \u00eenv\u0103luind locurile de apelare a serviciilor externe \u00een abstrac\u021bii corespunz\u0103toare. De asemenea, devine posibil\u0103 alegerea unei strategii pentru <noindex><a rel=\"nofollow\" href=\"https:\/\/docs.microsoft.com\/en-us\/azure\/architecture\/patterns\/bulkhead\">bulkhead<\/a><\/noindex> pattern, ceea ce poate fi de asemenea util. De exemplu, \u00een spring-cloud-netflix, acesta poate fi un thread pool sau un semafor.<\/p>\n<p><\/p>\n<h2 id=\"vyvod\">Ie\u0219ire<\/h2>\n<p><\/p>\n<p>Ca rezultat, am ob\u021binut o aplica\u021bie separat\u0103 care permite repetarea prelucr\u0103rii mesajului \u00een cazul \u00een care o sistem extern devine temporar indisponibil.<\/p>\n<p><\/p>\n<p>Unul dintre principalele avantaje ale aplica\u021biei este c\u0103 poate fi utilizat\u0103 de sistemele externe care func\u021bioneaz\u0103 pe acela\u0219i cluster Kafka, f\u0103r\u0103 modific\u0103ri semnificative de partea lor! Acelei aplica\u021bii \u00eei va fi necesar doar s\u0103 ob\u021bin\u0103 acces la topicul de retry, s\u0103 completeze c\u00e2teva antete Kafka \u0219i s\u0103 trimit\u0103 mesajul \u00een Retryer. Nu trebuie s\u0103 ridice nicio infrastructur\u0103 suplimentar\u0103. \u0218i, pentru a reduce num\u0103rul de mesaje transferate din aplica\u021bie \u00een Retryer \u0219i \u00eenapoi, am separat aplica\u021biile cu logic\u0103 liniar\u0103 \u0219i am implementat procesarea repetat\u0103 prin oprirea Consumer-ului.<\/p>\n<p>Sursa: <a content=\"nofollow\" rel=\"nofollow\" href=\"https:\/\/habr.com\/ru\/company\/tinkoff\/blog\/487094\/\">habr.com<\/a><\/p>","protected":false,"gt_translate_keys":[{"key":"rendered","format":"html"}]},"excerpt":{"rendered":"<p>\u041f\u0440\u0438\u0432\u0435\u0442, \u0425\u0430\u0431\u0440. \u041d\u0435\u0434\u0430\u0432\u043d\u043e \u044f \u043f\u043e\u0434\u0435\u043b\u0438\u043b\u0441\u044f \u043e\u043f\u044b\u0442\u043e\u043c \u043e \u0442\u043e\u043c, \u043a\u0430\u043a\u0438\u0435 \u043f\u0430\u0440\u0430\u043c\u0435\u0442\u0440\u044b \u043c\u044b \u0432 \u043a\u043e\u043c\u0430\u043d\u0434\u0435 \u0447\u0430\u0449\u0435 \u0432\u0441\u0435\u0433\u043e \u0438\u0441\u043f\u043e\u043b\u044c\u0437\u0443\u0435\u043c \u0434\u043b\u044f Kafka Producer \u0438 Consumer, \u0447\u0442\u043e\u0431\u044b \u043f\u0440\u0438\u0431\u043b\u0438\u0437\u0438\u0442\u044c\u0441\u044f \u043a \u0433\u0430\u0440\u0430\u043d\u0442\u0438\u0440\u043e\u0432\u0430\u043d\u043d\u043e\u0439 \u0434\u043e\u0441\u0442\u0430\u0432\u043a\u0435. \u0412 \u044d\u0442\u043e\u0439 \u0441\u0442\u0430\u0442\u044c\u0435 \u0445\u043e\u0447\u0443 \u0440\u0430\u0441\u0441\u043a\u0430\u0437\u0430\u0442\u044c, \u043a\u0430\u043a \u043c\u044b \u043e\u0440\u0433\u0430\u043d\u0438\u0437\u043e\u0432\u0430\u043b\u0438 \u043f\u043e\u0432\u0442\u043e\u0440\u043d\u0443\u044e \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0443 \u0441\u043e\u0431\u044b\u0442\u0438\u044f, \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u043d\u043e\u0433\u043e \u0438\u0437 Kafka, \u0432 \u0440\u0435\u0437\u0443\u043b\u044c\u0442\u0430\u0442\u0435 \u0432\u0440\u0435\u043c\u0435\u043d\u043d\u043e\u0439 \u043d\u0435\u0434\u043e\u0441\u0442\u0443\u043f\u043d\u043e\u0441\u0442\u0438 \u0432\u043d\u0435\u0448\u043d\u0435\u0439 \u0441\u0438\u0441\u0442\u0435\u043c\u044b. \u0421\u043e\u0432\u0440\u0435\u043c\u0435\u043d\u043d\u044b\u0435 \u043f\u0440\u0438\u043b\u043e\u0436\u0435\u043d\u0438\u044f \u0440\u0430\u0431\u043e\u0442\u0430\u044e\u0442 \u0432 \u043e\u0447\u0435\u043d\u044c \u0441\u043b\u043e\u0436\u043d\u043e\u0439 \u0441\u0440\u0435\u0434\u0435. \u0411\u0438\u0437\u043d\u0435\u0441-\u043b\u043e\u0433\u0438\u043a\u0430, \u043e\u0431\u0435\u0440\u043d\u0443\u0442\u0430\u044f [&hellip;]<\/p>\n","protected":false,"gt_translate_keys":[{"key":"rendered","format":"html"}]},"author":1,"featured_media":0,"comment_status":"open","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[688],"tags":[],"class_list":["post-56186","post","type-post","status-publish","format-standard","hentry","category-administrirovanie"],"aioseo_notices":[],"aioseo_head":"\n\t\t<!-- All in One SEO 5.0.2.1 - aioseo.com -->\n\t<meta name=\"description\" content=\"\u041f\u0440\u0438\u0432\u0435\u0442, \u0425\u0430\u0431\u0440. \u041d\u0435\u0434\u0430\u0432\u043d\u043e.\" \/>\n\t<meta name=\"robots\" content=\"max-image-preview:large\" \/>\n\t<meta name=\"author\" content=\"Yuri Gagarin\"\/>\n\t<link rel=\"canonical\" href=\"https:\/\/prohoster.info\/ro\/blog\/administrirovanie\/povtornaya-obrabotka-sobytij-poluchennyh-iz-kafka\" \/>\n\t<meta name=\"generator\" content=\"All in One SEO (AIOSEO) 5.0.2.1\" \/>\n\t\t<meta property=\"og:locale\" content=\"ro_RO\" \/>\n\t\t<meta property=\"og:site_name\" content=\"ProHoster | \u041a\u0443\u043f\u0438\u0442\u044c \u043d\u0430\u0434\u0435\u0436\u043d\u044b\u0439 \u0445\u043e\u0441\u0442\u0438\u043d\u0433 \u0434\u043b\u044f \u0441\u0430\u0439\u0442\u043e\u0432 \u0441 \u0437\u0430\u0449\u0438\u0442\u043e\u0439 \u043e\u0442 DDoS, VPS VDS \u0441\u0435\u0440\u0432\u0435\u0440\u044b\" \/>\n\t\t<meta property=\"og:type\" content=\"article\" \/>\n\t\t<meta property=\"og:title\" content=\"\ud83e\udd47\u041f\u043e\u0432\u0442\u043e\u0440\u043d\u0430\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0430 \u0441\u043e\u0431\u044b\u0442\u0438\u0439, \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u043d\u044b\u0445 \u0438\u0437 Kafka | ProHoster\" \/>\n\t\t<meta property=\"og:description\" content=\"\u041f\u0440\u0438\u0432\u0435\u0442, \u0425\u0430\u0431\u0440. \u041d\u0435\u0434\u0430\u0432\u043d\u043e.\" \/>\n\t\t<meta property=\"og:url\" content=\"https:\/\/prohoster.info\/ro\/blog\/administrirovanie\/povtornaya-obrabotka-sobytij-poluchennyh-iz-kafka\" \/>\n\t\t<meta property=\"og:image\" content=\"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg\" \/>\n\t\t<meta property=\"og:image:secure_url\" content=\"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg\" \/>\n\t\t<meta property=\"og:image:width\" content=\"350\" \/>\n\t\t<meta property=\"og:image:height\" content=\"350\" \/>\n\t\t<meta property=\"article:published_time\" content=\"2020-02-05T21:00:00+00:00\" \/>\n\t\t<meta property=\"article:modified_time\" content=\"2020-02-18T11:04:24+00:00\" \/>\n\t\t<meta property=\"article:publisher\" content=\"https:\/\/www.facebook.com\/prohoster\" \/>\n\t\t<meta property=\"article:author\" content=\"https:\/\/www.facebook.com\/prohoster\" \/>\n\t\t<!-- All in One SEO -->\n\n","aioseo_head_json":{"title":"\ud83e\udd47Reprelucrarea evenimentelor primite din Kafka | ProHoster","description":"Salut, Habr. Recent.","canonical_url":"https:\/\/prohoster.info\/ro\/blog\/administrirovanie\/povtornaya-obrabotka-sobytij-poluchennyh-iz-kafka","robots":"max-image-preview:large","keywords":"","webmasterTools":{"miscellaneous":""},"schema":null,"og:locale":"ro_RO","og:site_name":"ProHoster | \u041a\u0443\u043f\u0438\u0442\u044c \u043d\u0430\u0434\u0435\u0436\u043d\u044b\u0439 \u0445\u043e\u0441\u0442\u0438\u043d\u0433 \u0434\u043b\u044f \u0441\u0430\u0439\u0442\u043e\u0432 \u0441 \u0437\u0430\u0449\u0438\u0442\u043e\u0439 \u043e\u0442 DDoS, VPS VDS \u0441\u0435\u0440\u0432\u0435\u0440\u044b","og:type":"article","og:title":"\ud83e\udd47\u041f\u043e\u0432\u0442\u043e\u0440\u043d\u0430\u044f \u043e\u0431\u0440\u0430\u0431\u043e\u0442\u043a\u0430 \u0441\u043e\u0431\u044b\u0442\u0438\u0439, \u043f\u043e\u043b\u0443\u0447\u0435\u043d\u043d\u044b\u0445 \u0438\u0437 Kafka | ProHoster","og:description":"\u041f\u0440\u0438\u0432\u0435\u0442, \u0425\u0430\u0431\u0440. \u041d\u0435\u0434\u0430\u0432\u043d\u043e.","og:url":"https:\/\/prohoster.info\/ro\/blog\/administrirovanie\/povtornaya-obrabotka-sobytij-poluchennyh-iz-kafka","og:image":"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg","og:image:secure_url":"https:\/\/prohoster.info\/wp-content\/uploads\/2021\/11\/logo-350.jpg","og:image:width":350,"og:image:height":350,"article:published_time":"2020-02-05T21:00:00+00:00","article:modified_time":"2020-02-18T11:04:24+00:00","article:publisher":"https:\/\/www.facebook.com\/prohoster","article:author":"https:\/\/www.facebook.com\/prohoster"},"aioseo_meta_data":{"post_id":"56186","title":null,"description":null,"keywords":null,"keyphrases":null,"primary_term":null,"canonical_url":null,"og_title":null,"og_description":null,"og_object_type":"default","og_image_type":"default","og_image_url":null,"og_image_width":null,"og_image_height":null,"og_image_custom_url":null,"og_image_custom_fields":null,"og_video":null,"og_custom_url":null,"og_article_section":null,"og_article_tags":null,"twitter_use_og":false,"twitter_card":"default","twitter_image_type":"default","twitter_image_url":null,"twitter_image_custom_url":null,"twitter_image_custom_fields":null,"twitter_title":null,"twitter_description":null,"schema":{"blockGraphs":[],"customGraphs":[],"default":{"data":{"Article":[],"Course":[],"Dataset":[],"FAQPage":[],"Movie":[],"Person":[],"Product":[],"ProductReview":[],"Car":[],"Recipe":[],"Service":[],"SoftwareApplication":[],"WebPage":[]},"graphName":"","isEnabled":true},"graphs":[]},"schema_type":null,"schema_type_options":null,"pillar_content":false,"robots_default":true,"robots_noindex":false,"robots_noarchive":false,"robots_nosnippet":false,"robots_nofollow":false,"robots_noimageindex":false,"robots_noodp":false,"robots_notranslate":false,"robots_max_snippet":null,"robots_max_videopreview":null,"robots_max_imagepreview":"large","priority":null,"frequency":null,"local_seo":null,"seo_analyzer_scan_date":null,"breadcrumb_settings":null,"limit_modified_date":false,"reviewed_by":null,"ai":null,"created":"2021-02-28 19:29:19","updated":"2022-10-02 15:56:10","focus_keyword":null,"additional_keywords":null,"truseo_locale":null},"gt_translate_keys":[{"key":"link","format":"url"}],"_links":{"self":[{"href":"https:\/\/prohoster.info\/ro\/wp-json\/wp\/v2\/posts\/56186","targetHints":{"allow":["GET"]}}],"collection":[{"href":"https:\/\/prohoster.info\/ro\/wp-json\/wp\/v2\/posts"}],"about":[{"href":"https:\/\/prohoster.info\/ro\/wp-json\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"https:\/\/prohoster.info\/ro\/wp-json\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"https:\/\/prohoster.info\/ro\/wp-json\/wp\/v2\/comments?post=56186"}],"version-history":[{"count":0,"href":"https:\/\/prohoster.info\/ro\/wp-json\/wp\/v2\/posts\/56186\/revisions"}],"wp:attachment":[{"href":"https:\/\/prohoster.info\/ro\/wp-json\/wp\/v2\/media?parent=56186"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"https:\/\/prohoster.info\/ro\/wp-json\/wp\/v2\/categories?post=56186"},{"taxonomy":"post_tag","embeddable":true,"href":"https:\/\/prohoster.info\/ro\/wp-json\/wp\/v2\/tags?post=56186"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}