
Cum am ajuns în această situație?
Recent, a trebuit să lucrez la backendul unui proiect cu încărcare mare, în care era necesar să organizez execuția regulată a unui număr mare de sarcini de fundal, cu calcule complexe și cereri către servicii externe. Proiectul este asincron și, înainte să ajung eu, avea un mecanism simplu de cron pentru lansarea sarcinilor: un ciclu cu verificarea timpului curent și lansarea grupurilor de corutine prin gather — această abordare a fost acceptabilă până când numărul de corutine a depășit câteva zeci și sute, însă, când au ajuns la peste două mii, a fost necesar să mă gândesc la organizarea unei cozi adecvate de sarcini cu un broker, mai mulți lucrători și altele.
La început, am decis să încerc Celery, pe care l-am folosit anterior. Având în vedere natura asincronă a proiectului, m-am aprofundat în problemă și am văzut , precum și , creat de autorul articolului.
Voi spune așa, proiectul este foarte interesant și funcționează destul de bine în alte aplicații ale echipei noastre, iar autorul afirmă că a reușit să-l implementeze în producție, folosind un pool asincron. Dar, din păcate, nu mi s-a potrivit foarte bine, deoarece am descoperit probleme cu lansarea grupată a sarcinilor (vezi. ). La momentul redactării articolului era deja închis, totuși, activitatea a fost desfășurată pe parcursul unei luni. Oricum, mult succes autorului și toate cele bune, deoarece există deja lucruri funcționale în librărie… în general, eu sunt de vină și pentru mine instrumentul s-a dovedit a fi prea greu de utilizat. În plus, în unele sarcini erau câte 2-3 cereri http către servicii diferite, astfel că, chiar și cu optimizarea sarcinilor, creăm 4.000 de conexiuni tcp, aproximativ la fiecare 2 ore — nu este foarte bine… Mi-ar plăcea să creez o sesiune pentru un tip de sarcină la lansarea lucrătorilor. Mai multe detalii despre numărul mare de cereri prin aiohttp .
În acest context, am început să caut alternative și am găsit! Creatorii Celery, mai precis, cum am înțeles eu , au creat , inițial pentru proiectul . Faust este scrisă sub influența Kafka Streams și funcționează cu Kafka ca broker, de asemenea, pentru stocarea rezultatelor de la lucrători se folosește rocksdb, iar cel mai important este că biblioteca este asincronă.
De asemenea, poți arunca o privire celery și faust de la creatorii celui mai recent: diferențele dintre ele, diferențele brokerilor, realizarea unei sarcini de bază. Totul este destul de simplu, dar în faust atrage atenția o caracteristică plăcută — datele tipizate pentru transmiterea în topic.
Ce vom face?
Așadar, în această mică serie de articole, voi arăta cum să adunăm date în sarcini de fundal folosind Faust. Sursa pentru proiectul nostru exemplu va fi, după cum sugerează numele, . Voi demonstra cum să scriu agenți (sink, topicuri, partiții), cum să fac execuții regulate (cron), cele mai convenabile comenzi cli faust (un wrapper peste click), clustering simplu, și la final vom integra datadog (funcționând din cutie) și vom încerca să vedem ceva. Pentru stocarea datelor adunate, vom folosi mongodb și motor pentru conectare.
P.S. Judecând după încrederea cu care este scris punctul despre monitorizare, cred că cititorul, la sfârșitul ultimei articole, va arăta cam așa:

Cerinte pentru proiect
Având în vedere că deja am promis câteva lucruri, să întocmim o mică listă cu ceea ce ar trebui să poată serviciul:
- Să exporte valori mobiliare și o prezentare generală a acestora (inclusiv câștiguri și pierderi, bilanț, cash flow — pentru ultimul an) — regulat
- Să exporte date istorice (pentru fiecare an de tranzacționare, să găsească extremele prețului de închidere a tranzacțiilor) — regulat
- Să exporte ultimele date de tranzacționare — regulat
- Să exporte o listă configurată de indicatori pentru fiecare valoare mobiliare — regulat
După cum este de așteptat, alegem un nume proiectului la întâmplare: horton
Pregătim infrastructura
Titlul este, bineînțeles, puternic, totuși, tot ce trebuie să facem este să scriem o mică configurație pentru docker-compose cu kafka (și zookeeper — în același container), kafdrop (dacă vrem să vedem mesajele din topicuri), mongodb. Obținem [docker-compose.yml]() de următorul tip:
versiune: '3'
servicii:
db:
nume_container: horton-mongodb-local
imagine: mongo:4.2-bionic
comandă: mongod --port 20017
restart: întotdeauna
porturi:
- 20017:20017
mediu:
- MONGO_INITDB_DATABASE=horton
- MONGO_INITDB_ROOT_USERNAME=admin
- MONGO_INITDB_ROOT_PASSWORD=admin_password
kafka-service:
nume_container: horton-kafka-local
imagine: obsidiandynamics/kafka
restart: întotdeauna
porturi:
- "2181:2181"
- "9092:9092"
mediu:
KAFKA_LISTENERS: "INTERNAL://:29092,EXTERNAL://:9092"
KAFKA_ADVERTISED_LISTENERS: "INTERNAL://kafka-service:29092,EXTERNAL://localhost:9092"
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: "INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT"
KAFKA_INTER_BROKER_LISTENER_NAME: "INTERNAL"
KAFKA_ZOOKEEPER_SESSION_TIMEOUT: "6000"
KAFKA_RESTART_ATTEMPTS: "10"
KAFKA_RESTART_DELAY: "5"
ZOOKEEPER_AUTOPURGE_PURGE_INTERVAL: "0"
kafdrop:
nume_container: horton-kafdrop-local
imagine: 'obsidiandynamics/kafdrop:latest'
restart: întotdeauna
porturi:
- '9000:9000'
mediu:
KAFKA_BROKERCONNECT: kafka-service:29092
depinde_de:
- kafka-serviceNu este nimic complicat aici. Pentru Kafka, au fost declarați doi listeneri: unul (intern) pentru utilizare în rețeaua compusă și al doilea (extern) pentru cereri din exterior, de aceea l-am expus afară. 2181 este portul pentru Zookeeper. În privința celorlalte, cred că este clar.
Pregătim scheletul proiectului
În varianta de bază, structura proiectului nostru ar trebui să arate astfel:
horton
├── docker-compose.yml
└── horton
├── agents.py *
├── alphavantage.py *
├── app.py *
├── config.py
├── database
│ ├── connect.py
│ ├── cruds
│ │ ├── base.py
│ │ ├── __init__.py
│ │ └── security.py *
│ └── __init__.py
├── __init__.py
├── records.py *
└── tasks.py **Tot ce am marcat încă nu atingem, ci doar creăm fișiere goale.**
Am creat structura. Acum să adăugăm dependențele necesare, să scriem configurația și să ne conectăm la MongoDB. Nu voi prezenta textul complet al fișierelor în articol pentru a nu vă plictisi, ci voi face linkuri către versiunile necesare.
Să începeți cu dependențele și metadatele despre proiect —
Apoi, lansăm instalarea dependențelor și crearea virtualenv (sau, puteți crea singuri un folder venv și activați mediu):
pip3 install poetry (dacă nu este deja instalat)
poetry installAcum să creăm — credențiale și unde să ne conectăm. Imediat acolo putem plasa și datele pentru AlphaVantage. Și acum trecem la — extragem datele pentru aplicație din configurația noastră. Da, mărturisesc, am folosit biblioteca mea — .
În ceea ce privește conectarea la MongoDB — totul este foarte simplu. Am declarat pentru conectare și pentru CRUD-uri, astfel încât să facem mai ușor cereri pe colecții.
Ce va urma?
Articolul nu este foarte lung, deoarece aici vorbesc doar despre motivație și pregătire, așa că nu-mi supărați — promit că în partea următoare va fi acțiune și grafică.
Așadar, în această următoare parte vom:
- Scrie un mic client pentru alphavantage pe aiohttp cu cereri către endpoint-urile necesare.
- Vom crea un agent care va colecta date despre valori mobiliare și prețurile istorice ale acestora.
Sursa: habr.com
