
Bună, mă numesc Alexandru, lucrez la CIAN ca inginer și mă ocup de administrarea sistemelor și automatizarea proceselor infrastructurale. În comentariile uneia dintre articolele anterioare, am fost rugați să explicăm de unde obținem 4 TB de jurnale pe zi și ce facem cu ele. Da, avem multe jurnale și pentru procesarea lor am creat un cluster infrastructural separat, care ne permite să rezolvăm rapid problemele. În acest articol voi povesti despre cum am adaptat acest sistem pe parcursul unui an pentru a lucra cu fluxul de date în continuă creștere.
De unde am început

În ultimii câțiva ani, sarcina pe cian.ru a crescut foarte rapid și, până în trimestrul trei din 2018, numărul de vizitatori unici a atins 11,2 milioane pe lună. La acea vreme, în momente critice, pierdeam până la 40% din jurnale, ceea ce ne îngreuna să analizăm rapid incidentele și cheltuiam foarte mult timp și efort pentru a le rezolva. De asemenea, adesea nu reușeam să găsim cauza problemelor, iar acestea se repetau după o vreme. Era un coșmar, ceva trebuia să facem.
În acel moment, pentru stocarea jurnalelor, foloseam un cluster format din 10 noduri de date cu ElasticSearch versiunea 5.5.2 și setări standard pentru indecși. Acesta a fost implementat cu mai bine de un an în urmă ca o soluție populară și accesibilă: atunci fluxul de jurnale nu era foarte mare, nu avea rost să venim cu configurații non-standard.
Procesarea jurnalelor de intrare era asigurată de Logstash pe diverse porturi pe cinci coordonatori ElasticSearch. Un index, indiferent de dimensiune, constă din cinci șarduri. A fost organizată o rotație orară și zilnică; astfel, în fiecare oră, în cluster apăreau aproximativ 100 de șarduri noi. Atâta timp cât numărul jurnalelor nu era foarte mare, clusterul făcea față, iar nimeni nu punea accent pe setările sale.
Problemele creșterii rapide
Volumul jurnalelor generate a crescut foarte rapid, deoarece s-au suprapus două procese. Pe de o parte, numărul utilizatorilor serviciului creștea constant. Pe de altă parte, am început să trecem activ la o arhitectură bazată pe microservicii, divizând vechile noastre monolith-uri scrise în C# și Python. Câteva zeci de microservicii noi, care înlocuiau părți ale monolith-ului, generau cu mult mai multe jurnale pentru clusterul infrastructural.
Scalarea a dus la faptul că clusterul a devenit practic incontrolabil. Când logurile au început să sosească cu o viteză de 20.000 de mesaje pe secundă, rotația frecventă nefolositoare a crescut numărul de shard-uri până la 6.000, iar pe un nod erau mai mult de 600 de shard-uri.
Aceasta a dus la probleme cu alocarea memoriei, iar când un nod cădea, toate shard-urile începeau să migreze simultan, multiplicând traficul și încărcând celelalte noduri, ceea ce făcea practic imposibilă scrierea de date în cluster. În această perioadă, am rămas fără loguri. Iar în cazul unei probleme cu serverul pierdeam 1/10 din cluster în principiu. Cantitatea mare de indecși mici a adăugat complexitate.
Fără loguri, nu înțelegeam cauzele incidentului și riscam să călcăm pe aceleași greble din nou, iar în ideologia echipei noastre aceasta era inacceptabil, deoarece toate mecanismele noastre de lucru erau concepute exact pentru opus - să nu repetăm aceleași probleme. Așadar, aveam nevoie de un volum complet de loguri și de livrarea lor aproape în timp real, deoarece echipa de ingineri de gardă monitoriza alertele nu doar din metrici, ci și din loguri. Pentru a înțelege amploarea problemei - la acea vreme, volumul total de loguri era de aproximativ 2 TB pe zi.
Am stabilit obiectivul - să excluzem complet pierderea logurilor și să reducem timpul de livrare în clusterul ELK la maximum 15 minute în timpul urgențelor (această cifră a devenit ulterior KPI intern).
Mecanism nou de rotație și noduri hot-warm

Transformarea cluster-ului a început cu actualizarea versiunii ElasticSearch de la 5.5.2 la 6.4.3. Cluster-ul nostru de versiune 5 s-a prăbușit din nou, iar noi am decis să-l închidem și să-l actualizăm complet - oricum nu aveam loguri. Așa că această tranziție am realizat-o în doar câteva ore.
Cea mai amplă transformare în acest stadiu a fost implementarea pe trei noduri cu un coordonator ca buffer intermediar pentru Apache Kafka. Brokerul de mesaje ne-a scăpat de pierderile de loguri în timpul problemelor cu ElasticSearch. În același timp, am adăugat în cluster 2 noduri și am trecut la arhitectura hot-warm cu trei noduri „călduțe”, plasate în rack-uri diferite din data center. Pe acestea am redirecționat logurile care nu trebuiau pierdute sub nicio formă — nginx, precum și logurile de erori ale aplicațiilor. Pe celelalte noduri se trimiteau loguri minore — debug, warning etc., iar după 24 de ore migrăm logurile „importante” de pe nodurile „călduțe”.
Pentru a nu crește numărul de indecși de dimensiuni mici, am trecut de la rotația pe timp la mecanismul rollover. Pe forumuri erau multe informații despre faptul că rotația pe dimensiunea indexului este foarte nesigură, așa că am decis să folosim rotația pe numărul de documente din index. Am analizat fiecare index și am fixat numărul de documente după care ar trebui să se declanșeze rotația. Astfel, am obținut o dimensiune optimă a shardului — nu mai mult de 50 GB.
Optimizarea clusterului

Cu toate acestea, nu ne-am eliberat complet de probleme. Din păcate, apăreau totuși indecși mici: aceștia nu atingeau volumul stabilit, nu erau rotați și erau eliminați printr-o curățare globală a indecșilor mai vechi de trei zile, având în vedere că am eliminat rotația pe dată. Aceasta ducea la pierderi de date deoarece indexul dispare complet din cluster, iar încercarea de a scrie într-un index inexistent deteriora logica curator-ului pe care îl foloseam pentru gestionare. Alias-ul pentru scriere se transforma în index și strica logica rollover-ului, provocând o creștere necontrolată a unor indecși până la 600 GB.
De exemplu, pentru configurația de rotație:
curator-elk-rollover.yaml
---
actions:
1:
action: rollover
options:
name: "nginx_write"
conditions:
max_docs: 100000000
2:
action: rollover
options:
name: "python_error_write"
conditions:
max_docs: 10000000
În absența alias-ului rollover apărea o eroare:
ERROR alias "nginx_write" not found.
ERROR Failed to complete action: rollover. <type 'exceptions.ValueError'>: Unable to perform index rollover with alias "nginx_write".
Am lăsat soluționarea acestei probleme pentru următoarea iterație și ne-am ocupat de o altă întrebare: am trecut la logica de lucru pull a Logstash, care se ocupă cu procesarea jurnalelor de intrare (eliminarea informațiilor inutile și îmbogățirea acestora). L-am plasat în Docker, pe care îl lansăm prin docker-compose, unde am plasat și logstash-exporter, care oferă metrici în Prometheus pentru monitorizarea în timp real a fluxului de jurnale. Astfel, ne-am dat posibilitatea de a schimba în mod fluid numărul de instanțe Logstash responsabile pentru procesarea fiecărui tip de jurnale.
În timp ce ne perfectionam clusterul, traffic-ul pe cian.ru a crescut la 12,8 milioane de utilizatori unici pe lună. Ca rezultat, s-a întâmplat că transformările noastre nu țineau pasul cu modificările din producție, și ne-am confruntat cu faptul că nodurile "călduțe" nu făceau față sarcinilor și împiedicau livrarea întregii jurnale. Datele "fierbinți" le-am obținut fără probleme, dar pentru livrarea celorlalte a fost nevoie să intervenim manual și să facem rollover manual, pentru a distribui uniform indexii.
În același timp, scalarea și modificarea setărilor instanțelor Logstash din cluster a fost complicată de faptul că acesta era un docker-compose local, iar toate acțiunile se efectau manual (pentru a adăuga noi capete, era necesar să trecem manual prin toate serverele și să executăm docker-compose up -d peste tot).
Redistribuirea jurnalele
În septembrie acestui an, am continuat să despărțim monolitul, sarcina pe cluster creștea, iar fluxul de jurnale se apropia de 30 de mii de mesaje pe secundă.

Următoarea iterație am început-o prin actualizarea echipamentului. De la cinci coordonatori am trecut la trei, am înlocuit nodurile de date și am câștigat atât la bani, cât și la capacitatea de stocare. Pentru noduri folosim două configurații:
- Pentru nodurile "fierbinți": E3-1270 v6 / 960Gb SSD / 32 Gb x 3 x 2 (3 pentru Hot1 și 3 pentru Hot2).
- Pentru nodurile "călduțe": E3-1230 v6 / 4Tb SSD / 32 Gb x 4.
În această iterație am mutat indexul cu jurnalele de acces ale microserviciilor, care ocupă la fel de mult spațiu cât jurnalele nginx, în al doilea grup din cele trei noduri "fierbinți". Datele de pe nodurile "fierbinți" le păstrăm acum timp de 20 de ore, după care le transferăm pe cele "călduțe" împreună cu celelalte jurnale.
Am rezolvat problema dispariției indicilor mici prin reconfigurarea rotației acestora. Acum, indicii sunt rotați la fiecare 23 de ore, chiar și în cazul în care există puține date. Aceasta a crescut puțin numărul de sharduri (aproximativ 800), dar din punct de vedere al performanței clusterului, este acceptabil.
Astfel, în cluster s-au obținut șase noduri „fierbinți” și doar patru „călduțe”. Aceasta provoacă o mică întârziere în cereri pe intervale mari de timp, dar creșterea numărului de noduri în viitor va rezolva această problemă.
În această iterație, am remediat și problema lipsei scalării semi-automate. Pentru aceasta, am desfășurat un cluster infrastructural Nomad — similar cu cel deja desfășurat pe producție. Deocamdată, numărul de Logstash nu se modifică automat în funcție de încărcare, dar vom ajunge și la acest lucru.

Planuri de viitor
Configurația implementată se scalează excelent, iar acum stocăm 13,3 TB de date — toate jurnalele din ultimele 4 zile, necesare pentru analiza rapidă a alertelor. O parte din jurnale le transformăm în metrici, pe care le stocăm în Graphite. Pentru a facilita munca inginerilor, avem metrici pentru clusterul infrastructural și scripturi pentru repararea semi-automată a problemelor tipice. După creșterea numărului de noduri de date, planificată pentru anul următor, vom trece la stocarea datelor timp de 4 până la 7 zile. Aceasta va fi suficient pentru o muncă operativă eficientă, deoarece întotdeauna încercăm să investigăm incidentele cât mai repede posibil, iar pentru investigațiile pe termen lung avem datele de telemetrie.
În octombrie 2019, vizitările site-ului cian.ru au crescut la 15,3 milioane de utilizatori unici pe lună. Acesta a fost un test serios pentru soluția arhitecturală de livrare a jurnatelor.
Acum ne pregătim să actualizăm ElasticSearch la versiunea 7. Din păcate, va trebui să actualizăm mappingul multor indicii din ElasticSearch, deoarece acestea au fost migrate de la versiunea 5.5 și au fost declarate deprecated în versiunea 6 (în versiunea 7 pur și simplu nu există). Aceasta înseamnă că, în timpul actualizării, cu siguranță va apărea un forțat care ne va lăsa fără jurnale pentru o perioadă. Din versiunea 7 așteptăm în special Kibana cu o interfață îmbunătățită și filtre noi.
Am atins obiectivul principal: am încetat să pierdem jurnalele și am redus timpul de nefuncționare al clusterului de infrastructură de la 2-3 căderi pe săptămână la câteva ore de lucrări de service pe lună. Toată această muncă în producție este aproape imperceptibilă. Cu toate acestea, acum putem determina cu exactitate ce se întâmplă cu serviciul nostru, putem face rapid acest lucru în mod liniștit și fără îngrijorări că jurnalele se vor pierde. În general, suntem mulțumiți, fericiți și ne pregătim pentru noi realizări, despre care vom vorbi mai târziu.
Sursa: habr.com
