
Привет, Хабр.
Մեծ ընթացակարգեր, թե quais պեարբրեր մենք թիմում հաճախակի օգտագործում ենք Kafka Producer և Consumer-ի համար, որպեսզի հասնենք երաշխավորված մատակարարումների: Այս հոդվածում ցանկանում եմ պատմել, թե ինչպես ենք կազմակերպել իրադարձության կրկնահույսը, որը ստացվել է Kafka-ից, արտաքին համակարգի ժամանակավոր բացակայության պատճառով:
Ժամանակակից ծրագրերը աշխատում են շատ բարդ միջավայրում: Բիզնես-լոգիկան, որը պատված է ժամանակակից տեխնոլոգիական կույտով, աշխատելով Docker պատկերների մեջ, որոնք կառավարում է օրկեստրատոր, ինչպիսին է Kubernetes կամ OpenShift, և գործող մյուս ծրագրերի կամ ձեռնարկությունների լուծումների հետ հաղորդակցվում է ֆիզիկական և վիրտուալ ռաուտերների շղթայով: Այդպիսի միջավայրում ամեն ինչ կարող է ընդհատվել, ուստի արտաքին համակարգերից մեկի բացակայության դեպքում իրադարձությունների կրկնահուզումը մեր բիզնեսի գործընթացների կարևոր մասն է:
Ինչպես եղել է մինչև Kafka
Ավելի վաղ նախագծում մենք օգտագործում էինք IBM MQ հաղորդագրությունների ասինխրոն մատակարարումների համար: Եթե որևէ սխալ է տեղի ունեցել ծառայության աշխատանքում, ստացված հաղորդագրությունը կարող էր տեղադրվել մեռյալ-խնդիրների հերթում (DLQ) հետագա ձեռքի դասավորության համար: DLQ-ն ստեղծվում էր մուտքային հերթի կողքին, հաղորդագրության փոխադրումը տեղի էր ունենում IBM MQ-ի ներսում:
Եթե սխալը ժամանակավոր բնույթ է ունեցել և մենք կարող էինք դա որոշել (օրինակ, ResourceAccessException HTTP-կանգառի ժամանակ կամ MongoTimeoutException MongoDb-ի հարցման ժամանակ), ապա գործում էի կրկնակի զանգերի ռազմավարությունը: Անկախ ծրագրի տրամաբանության ճյուղից, սկզբնական հաղորդագրությունը տեղադրվում էր կամ համակարգային հերթի համար ուշացումով տեղադրման կամ տարբեր եղանակի համար, որը վաղուց պատրաստվել էր հաղորդագրությունների կրկնակի մատակարարումի համար: Այս պահին հաղորդագրության վերնագրում գրանցվում է կրկնակի մատակարարումների համարը, որը կապվում է ուշացման միջակայքի կամ ծրագրի մակարդակի ռազմավարության վերջի հետ: Եթե մենք հասնում ենք ռազմավարության վերջի, բայց արտաքին համակարգը դեռ բացակայում է, ապա հաղորդագրությունը տեղադրվում է DLQ-ում ձեռքի դասավորության համար:
Լուծման որոնում
, կարող եք գտնել հետևյալը: . Եթե հակիրճ, ապա առաջարկվում է մի թեմա ստեղծել յուրաքանչյուր ուշացման միջակայքի համար և իրականացնել հաղորդագրությունները կարդացողներ, որոնք միջոցառումներ կկատարեն անհրաժեշտ ժամանակով:

Դեռևս մեծ քանակությամբ դրական ակնարկների առկայության notwithstanding, այն մեզ չի թվայում համենայնիև, քանի որ մշակողը, բացի բիզնես պահանջների իրականացումից, երկար ժամանակ կներդնի նկարագրված մեխանիզմի իրագործման համար:
Բացի այդ, եթե Kafka կլաստերում գործողությունը մուտքի կառավարում է, ապա հարկավոր է որոշ ժամանակ ծախսել թեմաներ ստեղծելու և դրանց անհրաժեշտ մուտք ապահովելու վրա։ Վերաշավելու բոլոր թեմաների համար պետք է ընտել ճիշտ retention.ms պարամետրը, որպեսզի հաղորդագրությունները կարողանան կրկին ուղարկվել և չեն կորչեն։ Հասանելիությունը և մուտքը պետք է կրկնել յուրաքանչյուր գործող կամ նոր ծառայի համար:
Այժմ դիտենք, թե ինչպիսի մեխանիզմներ ներկայացնում են հաղորդագրությունների կրկնակի մշակում ծավալուն ձևի համար Spring-ում և մասնավորապես spring-kafka-ում։ Spring-kafka-ն ունի տրանսիտիվ կախվածություն spring-retry-ից, որը ապահովում է տարբեր BackOffPolicy-ների կառավարման աբստրակցիաներ։ Սա բավականին ճկուն գործիք է, բայց դրա նշանակալի բացասական կողմը հաղորդագրությունների պահպանումն է հայտնելու դիմաց application's memory-ում։ Սա նշանակում է, որ ծրագիրը վերագործարկելու դեպքում, օրինակ՝ թարմացումների կամ օգտագործման ընթացքում սխալի պատճառով, կկորսնվի բոլոր հաղորդագրությունները, որոնք սպասում էին կրկնակի մշակման։ Որպեսզի այս հարցը վճռական չէր մեր համակարգի համար, մենք չքննարկեցինք այն ավելի շատ։
Այնուամենայնիվ, spring-kafka-ն ներկայացնում է несколько реализаций ContainerAwareErrorHandler՝ օրինակ , որի միջոցով կարելի է, սխալի դեպքում offset-ը չտեղափոխելուց օգտվել և մշակել հաղորդագրությունը ավելի ուշ։ Եթե spring-kafka 2.3-ից սկսելուց, BackOffPolicy առաջարկելու հնարավորությունը появилась։
Այս մոտեցումը թույլ է տալիս կրկնակի մշարված հաղորդագրություններին ապրելու ծրագիրը վերարկվել հրատապի ժամանակ, բայց DLQ մեխանիզմը դեռևս բացակայում է։ Աճին, մենք ընտրել ենք այս տարբերակը 2019 թվականի սկզբին՝ օպտիմիստիկ կարծիքով կարծում էինք, որ DLQ չի լինի (մենք բարեհաջող էինք և այն իսկապես չէր понадобվեց application's эксплуатаций несколько ամիսներից հետո այս կրկնակի մշակման համակարգով)։ Տեղական սխալները անդրադարձնում էին SeekToCurrentErrorHandler-ի կրակմանը։ Մյուս սխալները տպվում էին տպագրության մեջ, առաջանում էին offset-ի տեղափոխման, և մշակումը շարունակվում էր հաջորդ հաղորդագրության հետ:
Վայելված լուծում
SeekToCurrentErrorHandler-ի հիման վրա իրականացումը մեզ տարանքի ենթարկեց սեփական մեխանիզմի մշակման՝ հաղորդագրությունները կրկին ուղարկելու համար։
Ավելի կարևորն այն է, որ ցանկանում էինք օգտագործել արդեն մերժված փորձը և ընդլայնել այն՝ ասյալռնալու տրամաբանությունից։ Լինար գործողություն ունենալու դեպքում, նպատակահարմար կլիներ ժամանակավոր ընդմիջում կանգնեցնելու նոր հաղորդագրությունների ընթերցումը՝ բաժանա վերակայման ռազմավարության շրջանակներում։ Մնացած ծրագրերի համար ցանկալի էր ունենալ ընդհանուր կետ, որը կդառնա կրկնորդում իրականացնող ռազմավարության կատարելու համակենտրոն։ Բացի այդ, այս ընդհանուր կետը պետք է ունենա DLQ ֆունկցիոնալություն երկու մոտեցումների համար։
Կրկնորդման ռազմավարությունն ինքը պետք է պահվի այն հավելվածում, որը պատասխանատվություն է կրում ժամանակավոր սխալի հայտնվելու դեպքում՝ հաջորդ միջվայրկյանի ստանալու համար։
Լինեի գործողությունը ծրագրվածության դեպքում, որը ունի մակարդակային տրամաբանություն
spring-kafka-ով աշխատելու ժամանակ, Consumer-ի կանգնեցնելու կոդը կարող է выглядеть примерно так:
public void pauseListenerContainer(MessageListenerContainer listenerContainer,
Instant retryAt) {
if (nonNull(retryAt) && listenerContainer.isRunning()) {
listenerContainer.stop();
taskScheduler.schedule(() -> listenerContainer.start(), retryAt);
return;
}
// DLQ-ի համար
}Օրինակ, retryAt-ի մեջ ժամանակն է, երբ MessageListenerContainer-ը պետք է նորից գործարկվի, եթե դեռ աշխատում է։ Նորից գործարկվելը տեղի կունենա TaskScheduler-ում, որի իրականացումը նույնպես ապահովում է spring-ը։
retryAt արժեքը ստանում ենք հետևյալ եղանակով՝
- Փնտրվում է կրկնորդման հաշվիչի արժեքը։
- Հաշվիչի արժեքին համապատասխան սկզբնավորվում է ծանրաբեռնման ներկայիս ժամանակահատվածը։ Ռազմավարությունը հայտարարվում է ինքնին հավելվածում, դրա պահման համար մենք ընտրել ենք JSON ձևաչափը։
- JSON-մասիվում գտնված ժամանակահատվածը պարունակում է այն վայրկյանների քանակը, որոնց ընթացքում պետք է կրկին մշակվի։ Այս վայրկյանները ավելացվում են ներկայիս ժամանակին՝ ստեղծելով retryAt արժեքը։
- Եթե ժամանակահատվածը չի գտնվել, ապա retryAt արժեքը հավասար է null-ին, և հաղորդագրությունը կուղարկվի DLQ՝ ձեռքով մշակման համար։
Այս մոտեցման շնորհիվ միայն պետք է պահպանել յուրաքանչյուր հաղորդագրության համար կրկնակի կանչերի քանակը, որը հիմա ընթացքի մեջ է, օրինակ՝ application's հիշողությանը: Այս մոտեցման համար փորձերի հաշվիչը հիշողության մեջ պահպանելը կարևոր չէ, քանի որ կիրառումը, որն ունի գծային տրամաբանություն, չի կարող համընդհանուր մշակման իրականացնել: Դասական spring-retry-ի հակադրությամբ, կիրառման վերակառուցումը չգիտի բոլոր հաղորդագրությունները վերամշակել, այլ просто стратегияן համար վերսկզբում:
Այս մոտեցումը օգնում է բեռը հեռացնե՛լ արտաքին համակարգից, որը կարող է չլինել հասանելի շատ մեծ բեռի պատճառով: Այլ կերպ ասած, կրկնակի վերամշակման կրկին հետ, մենք հասել ենք տարբերակի իրականացմանը, .
Մեր դեպքում սխալների շեմը ընդամենը 1 է, իսկ ժամանակավոր ցանցային ընդհատումից համակարգի դանդաղեցումը նվազագույնի հասցնելու համար, մենք օգտագործում ենք շատ մանրամասն կրկնակի կանչերի ռազմավարություն՝ փոքր սպասումների միջակայքներով: Սա կարող է չհամապատասխանել բոլոր ընկերության խմբի հավելվածներին, ուստի սխալների շեմի և միջակայքի չափը պետք է ընտրել համակարգի առանձնահատկությունները հիմք ունենալով:
Ապահովիչ հավելված հաղորդագրություններ վարելու համար՝ չհրատապ տրամաբանական տրամադրությամբ ազդանշաններ
Ահա մեկնաբանության օրինակ կոդը, որը փոխանցում է հաղորդագրություն նմանատիպ հավելվածում (Retryer), որն իրականացրել է կրկնակի փոխանցում DESTINATION թեմայով, երբ RETRY_AT ժամանակն հասնում է:
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);
} Օրինակով երևում է, որ շատ տեղեկություններ փոխանցվում են գլուխներում: RETRY_AT արժեքը գտնվում է այնպես, ինչպես նորից պատճենելու մեխանիզմում՝ դադարեցնող սպառողին: Ավելին, մենք փոխանցում ենք DESTINATION եւ RETRY_AT:
- GROUP_ID՝ ըստ որի խմբավորում ենք հաղորդագրությունները ձեռքով վերլուծության եւ որոնումն ավելի հեշտացնելու համար:
- ORIGINAL_PARTITION, որպեսզի փորձենք պահել նույն սպառողը կրկնակի մշակման համար: Այս պարամետրը կարող է հավասար լինել null, այս դեպքում նոր partition-ը կստեղծվի record.key()-ի սկզբնական հաղորդագրության բանալիով:
- Թարմացված COUNTER արժեք, որպեսզի հետեւենք կրկնակի զանգերի ռազմավարությանը:
- SEND_TO — տվյալ, որը ցույց է տալիս, թե հաղորդագրությունը պետք է կրկնակի մշակման ուղարկվի RETRY_AT հասնելուց հետո, թե հարկ է տեղադրել DLQ:
- REASON — պատճառը, թե ինչու հաղորդագրության մշակումը դադարեցվել է:
Retryer-ը պահպանում է հաղորդագրությունները կրկնակի ուղարկման եւ ձեռքով ուսումնասիրման համար PostgreSQL-ում: Ժամաչափով գործարկվում է առաջադրանք, որը գտնում է RETRY_AT ժամկետը անցած հաղորդագրությունները եւ ուղարկում դրանք ORIGINAL_PARTITION partition-ը DESTINATION թեմա, օգտագործելով record.key()-ը:
Հաղորդագրությունները ուղարկելուց հետո հեռացվում են PostgreSQL-ից: Հաղորդագրությունների ձեռքով ուսումնասիրման գործընթացը տեղի է ունենում պարզ UI-ում, որը համագործակցում է Retryer-ի հետ REST API-ի միջոցով: Նրա հիմնական հատկությունները ներառում են հաղորդագրությունների կրկնակի ուղարկում կամ հեռացում DLQ-ից, սխալների տեղեկությունների դիտում եւ հաղորդագրությունների որոնում, օրինակ՝ սխալի անվանով:
Քանի որ մեր կլաստերներում բացառված է մուտքի կառավարումը, անհրաժեշտ է հավելյալ մուտք կամացյալ ահազանգել ընթերցվող թեմայից, Դրա համար պահանջվում է, որ Retryer-ը կարողանա գրել DESTINATION թեմային: Սա անհարմար է, սակայն, ի տարբերություն ժամանակահատվածի թեմայի մոտեցման, մենք ստանում ենք լիարժեք DLQ եւ UI նրա կառավարելու համար:
Կան դեպքեր, երբ մուտքային թեման ընթերցում են մի քանի տարբեր սպառողային խմբեր, որոնց ծրագրերը իրականացնում են տարբեր տրամաբանություն: Մի ծրագրի համար Retryer-ի միջոցով հաղորդագրության կրկնակի մշակումը կհանգեցնի կրկնօրինակին другой-ին: Սա խուսափելու համար, մենք ստեղծում ենք առանձնահատուկ թեմա կրկնակի մշակման համար: Մուտքային եւ retry թեման կարող է կարդալ նույն սպառողը առանց որեւէ սահմանափակման:

Ավելորդ իմաստով այս մոտեցումը չի տրամադրում circuit breaker-ի հնարավորությունը, սակայն կարելի է այն ավելացնել ծրագրում՝ օգտագործելով կամ նոր , արտաքին ծառայությունների 호출ի վայրերը համապատասխան աբստրակցիաների մեջ վերածելով: Ավելին, վարչապետության ընտրության ռազմավարության հնարավորությունը։ պատառիկ, ինչը նույնպես կարող է օգտակար լինել։ Օրինակ, spring-cloud-netflix համակարգում դա կարող է լինել թելային հանրապետություն կամ սեմաֆոր։
Ամփոփում
Վերատվության արդյունքում ստացվել է առանձին ծրագիր, որը թույլ է տալիս կրկին իրականացնել հաղորդագրության մշակումը, երբ ժամանակավորապես անհասանելի է որևէ արտաքին համակարգ։
Ծրագրի գլխավոր առավելություններից մեկն այն է, որ այն կարող են օգտագործել արտաքին համակարգեր, որոնք աշխատում են նույն Kafka խմբի հետ, առանց մեծ փոփոխությունների իրենց կողմում։ Այս ծրագրին միայն անհրաժեշտ կլինի սույնին մուտք գործել retry-թեմա, լրացնել մի քանի Kafka վերնագրեր և ուղարկել հաղորդագրությունը Retryer-ին։ Ոչ մի լրացուցիչ ենթակառուցվածք բարձրացնել պետք չէ։ Իսկ հաղորդագրությունների քանակը, որը տեղափոխվում է ծրագրից Retryer և հակառակ ուղղությամբ, փոքրացնելու համար, առանձնացման ենք կատարել գիծ առաջատար տրամադրություն ունեցող ծրագրերը և իրականացնում ենք կրկնակի մշակումը՝ կանգնած Consumer-ի միջոցով։
Ընտանիք: habr.com
