
Hoe ben ik hier beland?
Niet zo lang geleden moest ik werken aan de backend van een zwaarbelast project, waarbij een regelmatige uitvoering van een groot aantal achtergrondtaken met complexe berekeningen en verzoeken naar externe services moest worden georganiseerd. Het project is asynchroon en voordat ik kwam, was er een eenvoudige cron-methode voor het starten van taken: een loop met een controle van de huidige tijd en het starten van groepen coroutines via gather — deze aanpak bleek acceptabel totdat het aantal coroutines in de tientallen en honderden begon te lopen, maar toen het aantal meer dan tweeduizend overschreed, moest ik gaan nadenken over een goede organisatie van een taakqueue met een broker, meerdere workers en andere zaken.
Eerst besloot ik Celery uit te proberen, wat ik eerder had gebruikt. Gezien de asynchroniciteit van het project, dook ik in het onderwerp en zag , evenals , gemaakt door de auteur van het artikel.
Ik zal het zo zeggen, het project is erg interessant en werkt behoorlijk succesvol in andere toepassingen van ons team, en de auteur zelf zegt ook dat hij het in productie heeft kunnen uitrollen, gebruikmakend van een asynchrone pool. Maar helaas paste het niet zo goed bij me, omdat er was met het groepsgewijs starten van taken (zie ). Op het moment dat dit artikel werd geschreven, was het al opgelost, echter, het werk liep een maand door. Hoe dan ook, veel succes en het allerbeste voor de auteur, omdat er al werkende dingen in de lib zijn... in ieder geval, het ligt aan mij en voor mij bleek de tool nog steeds niet helemaal klaar. Bovendien waren er in sommige taken 2-3 http-verzoeken naar verschillende services, hierdoor creëren we zelfs bij optimalisatie van taken 4000 tcp-verbindingen, ongeveer elke 2 uur — dat is niet echt ideaal... Het zou fijn zijn om een sessie voor één type taak op te zetten bij het starten van workers. Wat betreft het grote aantal verzoeken via aiohttp, .
daarom begon ik te zoeken naar alternatieven en ik vond het! De makers van Celery, en specifiek, zoals ik het begreep, , hebben , oorspronkelijk gemaakt voor het project . Faust is geschreven ter inspiratie van Kafka Streams en werkt met Kafka als broker, ook wordt rocksdb gebruikt voor de opslag van resultaten van de agenten, en het belangrijkste is dat de bibliotheek asynchroon is.
Daarnaast kun je ook kijken naar Celery en Faust van de makers van de nieuwste: hun verschillen, de verschillen tussen de brokers, de implementatie van een eenvoudige taak. Het is allemaal vrij simpel, maar in Faust valt een aangename functie op — getypeerde gegevens voor overdracht naar een topic.
Wat gaan we doen?
Dus, in een kleine serie artikelen laat ik zien hoe je gegevens kunt verzamelen in achtergrondtaken met behulp van Faust. De bron voor ons voorbeeldproject zal, zoals de titel al aangeeft, zijn. Ik zal demonstreren hoe je agenten (sink, topics, partitions) schrijft, hoe je regelmatig (cron) uitvoert, handige cli-commando's van Faust (een wrapper boven click), eenvoudige clustering, en aan het eind zullen we Datadog (direct werkend) aansluiten en proberen iets te zien. Voor het opslaan van de verzamelde gegevens zullen we MongoDB en Motor voor de connectie gebruiken.
P.S. Gezien het vertrouwen waarmee de sectie over monitoring is geschreven, denk ik dat de lezer aan het eind van het laatste artikel er ongeveer zo zal uitzien:

Vereisten voor het project
Aangezien ik al wat beloofd heb, laten we een korte lijst opstellen van wat de service moet kunnen:
- Waarde effecten en overzicht daarvan exporteren (inclusief winst en verlies, balans, cashflow — voor het afgelopen jaar) — regelmatig
- Historische gegevens exporteren (extremen van de sluitprijs voor elk handelsjaar vinden) — regelmatig
- Laatste handelsgegevens exporteren — regelmatig
- Een geconfigureerde lijst van indicatoren voor elk effect exporteren — regelmatig
Zoals het hoort, kiezen we een naam voor het project uit de lucht: horton
De infrastructuur voorbereiden
De titel is natuurlijk sterk, maar alles wat we moeten doen is een kleine configuratie voor docker-compose schrijven met Kafka (en Zookeeper — in één container), Kafdrop (als we de berichten in de topics willen bekijken), MongoDB. We krijgen [docker-compose.yml]() van de volgende aard:
versie: '3'
diensten:
db:
container_name: horton-mongodb-local
image: mongo:4.2-bionic
command: mongod --port 20017
restart: altijd
ports:
- 20017:20017
environment:
- MONGO_INITDB_DATABASE=horton
- MONGO_INITDB_ROOT_USERNAME=admin
- MONGO_INITDB_ROOT_PASSWORD=admin_password
kafka-service:
container_name: horton-kafka-local
image: obsidiandynamics/kafka
restart: altijd
ports:
- "2181:2181"
- "9092:9092"
environment:
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:
container_name: horton-kafdrop-local
image: 'obsidiandynamics/kafdrop:latest'
restart: altijd
ports:
- '9000:9000'
environment:
KAFKA_BROKERCONNECT: kafka-service:29092
depends_on:
- kafka-serviceDit is helemaal niet ingewikkeld. Voor kafka hebben we twee listeners gedefinieerd: één (intern) voor gebruik binnen het compositienetwerk en de andere (extern) voor verzoeken van buitenaf, daarom hebben we deze naar buiten doorgestuurd. 2181 is de poort van de zookeeper. Verder denk ik dat het duidelijk is.
Laten we de basis van het project creëren.
In de basisstructuur zou ons project er als volgt uit moeten zien:
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 **Alles wat ik heb gemarkeerd laten we voorlopig nog even met rust, terwijl we gewoon lege bestanden aanmaken.**
We hebben de structuur aangemaakt. Laten we nu de benodigde afhankelijkheden toevoegen, de configuratie schrijven en de verbinding met mongodb maken. Ik zal de volledige tekst van de bestanden niet in het artikel opnemen, om het niet te langdradig te maken, maar ik zal links maken naar de benodigde versies.
Laten we beginnen met de afhankelijkheden en meta-informatie over het project —
Daarna starten we de installatie van de afhankelijkheden en maken we de virtualenv aan (of u kunt zelf een map venv aanmaken en de omgeving activeren):
pip3 install poetry (als het nog niet is geïnstalleerd)
poetry installLaten we nu — de credentials en waar te verbinden. Hier kunnen we ook de gegevens voor alphavantage plaatsen. En dan gaan we verder met — we halen gegevens voor de applicatie uit onze configuratie. Ja, ik geef het toe, ik heb mijn eigen bibliotheek gebruikt — .
Wat betreft de verbinding met mongo — het is helemaal niet ingewikkeld. We hebben voor verbinding gedefinieerd en voor de CRUD operaties, zodat het gemakkelijker is om verzoeken naar de collecties te doen.
Wat zal er daarna gebeuren?
Het artikel is niet zo lang geworden, omdat ik hier alleen over motivatie en voorbereiding spreek, dus mijn excuses — ik beloof dat er in het volgende deel actie en graphics zullen komen.
Dus, in dit volgende deel gaan we:
- Een kleine client voor alphavantage schrijven met aiohttp die verzoeken naar de benodigde eindpunten verzendt.
- Een agent maken die gegevens over effecten verzamelt en historische prijzen daarvoor opzoekt.
Bron: habr.com
