
Wie bin ich zu diesem Leben gekommen?
KĂŒrzlich musste ich an einem stark ausgelasteten Backend-Projekt arbeiten, bei dem es notwendig war, regelmĂ€Ăig eine groĂe Anzahl von Hintergrundaufgaben mit komplexen Berechnungen und Anfragen an externe Dienste zu organisieren. Das Projekt ist asynchron, und bevor ich kam, gab es einen einfachen Mechanismus zum Kicken von Aufgaben: eine Schleife mit der ĂberprĂŒfung der aktuellen Zeit und dem Start von Gruppen von Koroutinen ĂŒber gather â dieser Ansatz war akzeptabel, solange die Anzahl der Koroutinen in den Zehner- und Hundertbereich fiel, jedoch als ihre Zahl ĂŒber zweitausend stieg, musste ich ĂŒber die Organisation einer normalen Aufgabenwarteschlange mit einem Broker, mehreren Workern und anderem nachdenken.
ZunÀchst habe ich beschlossen, Celery auszuprobieren, das ich zuvor verwendet hatte. Angesichts der AsynchronitÀt des Projekts tauchte ich in das Thema ein und sah , sowie , erstellt vom Autor des Artikels.
Ich sage mal so, das Projekt ist sehr interessant und funktioniert ziemlich gut in anderen Anwendungen unseres Teams. Der Autor selbst erwĂ€hnt, dass er es erfolgreich in Produktion gebracht hat, indem er einen asynchronen Pool genutzt hat. Leider war das jedoch nicht ganz passend fĂŒr mich, da sich herausstellte, mit dem gruppierten Start von Aufgaben (siehe ). Zum Zeitpunkt der Artikelverfassung bereits behoben, jedoch dauerte die Arbeit einen Monat. Auf jeden Fall wĂŒnsche ich dem Autor viel Erfolg und alles Gute, da es bereits funktionierende Lösungen in der Bibliothek gibt⊠im Grunde liegt es an mir, und das Tool erschien mir noch nicht ausgereift. AuĂerdem gab es in einigen Aufgaben 2-3 HTTP-Anfragen an verschiedene Dienste, wodurch wir selbst bei der Optimierung der Aufgaben 4000 TCP-Verbindungen alle zwei Stunden erzeugen â nicht ideal⊠Ich wĂŒrde gerne eine Sitzung fĂŒr einen Typ von Aufgaben beim Starten der Worker erzeugen. Ein wenig mehr ĂŒber die groĂe Anzahl von Anfragen ĂŒber aiohttp. .
In diesem Zusammenhang begann ich, nach Alternativen zu suchen und fand! Von den Erstellern von Celery, speziell, wie ich verstanden habe, , wurde geschaffen, ursprĂŒnglich fĂŒr das Projekt . Faust wurde von Kafka Streams inspiriert und verwendet Kafka als Broker. FĂŒr die Speicherung der Ergebnisse der Agenten setzen wir rocksdb ein. Das Wichtigste ist jedoch, dass die Bibliothek asynchron ist.
AuĂerdem können Sie einen Blick auf von Celery und Faust, erstellt von den Entwicklern von Faust, werfen: ihre Unterschiede, die Broker-Differenzen und die Implementierung einer grundlegenden Aufgabe. Alles ist ziemlich einfach, aber besonders interessant an Faust ist die Funktion der typisierten Daten, die an das Topic ĂŒbergeben werden.
Was werden wir tun?
In einer kleinen Artikelsammlung zeige ich, wie man mit Faust Daten in Hintergrundaufgaben sammelt. Die Quelle fĂŒr unser Beispielprojekt wird, wie der Titel schon sagt, . Ich werde demonstrieren, wie man Agenten (Sink, Topics, Partitionen) erstellt, regelmĂ€Ăige (Cron) AusfĂŒhrungen einrichtet, die nĂŒtzlichen CLI-Befehle von Faust (eine Wrapper ĂŒber Click) verwendet, einfaches Clustering durchfĂŒhrt und am Ende Datadog (aus der Box funktionierend) anbindet, um zu sehen, was wir erreichen können. FĂŒr die Speicherung der gesammelten Daten verwenden wir MongoDB und Motor fĂŒr die Verbindung.
P.S. Angesichts des Selbstbewusstseins, mit dem der Punkt ĂŒber das Monitoring geschrieben ist, denke ich, dass der Leser am Ende des letzten Artikels so aussehen wird:

Projektanforderungen
Da ich bereits einige Versprechungen gemacht habe, erstellen wir eine kleine Liste dessen, was der Service können sollte:
- RegelmĂ€Ăig wertpapierbezogene Daten und einen Ăberblick darĂŒber extrahieren (einschlieĂlich Gewinne und Verluste, Bilanz, Cashflow â fĂŒr das letzte Jahr)
- RegelmĂ€Ăig historische Daten extrahieren (fĂŒr jedes Handelsjahr die Extrempreise der Schlusskurse finden)
- RegelmĂ€Ăig die letzten Handelsdaten extrahieren
- RegelmĂ€Ăig eine konfigurierte Liste von Indikatoren fĂŒr jedes Wertpapier extrahieren
Wie es sich gehört, wĂ€hlen wir einen Projektnamen willkĂŒrlich aus: horton
Infrastruktur vorbereiten
Die Ăberschrift ist stark, aber alles, was wir tun mĂŒssen, ist, eine kleine Konfiguration fĂŒr docker-compose mit Kafka (und Zookeeper â in einem Container), Kafdrop (falls wir die Nachrichten in den Themen sehen möchten) und MongoDB zu schreiben. Wir erhalten [docker-compose.yml]() in folgender Form:
version: '3'
services:
db:
container_name: horton-mongodb-local
image: mongo:4.2-bionic
command: mongod --port 20017
restart: always
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: always
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: always
ports:
- '9000:9000'
environment:
KAFKA_BROKERCONNECT: kafka-service:29092
depends_on:
- kafka-serviceEs ist wirklich nichts kompliziert. FĂŒr Kafka wurden zwei Listener definiert: einer (intern) fĂŒr die Nutzung innerhalb des Composite-Netzwerks und der andere (extern) fĂŒr Anfragen von auĂen, weshalb er nach auĂen geleitet wurde. 2181 ist der Port des Zookeepers. DarĂŒber hinaus ist alles klar.
Wir bereiten das GrundgerĂŒst des Projekts vor.
Im Basisversion sollte die Struktur unseres Projekts wie folgt aussehen:
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, was ich markiert habe berĂŒhren wir vorerst nicht, sondern erstellen einfach leere Dateien.**
Wir haben die Struktur erstellt. Jetzt fĂŒgen wir die notwendigen AbhĂ€ngigkeiten hinzu, schreiben die Konfiguration und stellen die Verbindung zu MongoDB her. Die vollstĂ€ndigen Inhalte der Dateien werde ich in diesem Artikel nicht auflisten, um es nicht unnötig in die LĂ€nge zu ziehen, sondern ich werde Links zu den benötigten Versionen einfĂŒgen.
Beginnen wir mit den AbhĂ€ngigkeiten und den Metadaten zum Projekt â
Danach starten wir die Installation der AbhÀngigkeiten und die Erstellung eines virtualenv (alternativ können Sie auch selbst einen Ordner venv erstellen und die Umgebung aktivieren):
pip3 install poetry (falls noch nicht installiert)
poetry installJetzt erstellen wir â die Credentials und wo wir uns verbinden. Dort können auch die Daten fĂŒr Alphavantage hinterlegt werden. Nun gehen wir zu â wir extrahieren die Daten fĂŒr die Anwendung aus unserer Konfiguration. Ja, ich gebe es zu, ich habe meine eigene Bibliothek verwendet â .
Die Verbindung zu MongoDB ist ganz einfach. Wir haben fĂŒr die Verbindung und fĂŒr Skripte, um Abfragen zu Sammlungen zu erleichtern.
Wie geht es weiter?
Der Artikel ist nicht besonders lang, da ich hier nur ĂŒber Motivation und Vorbereitung spreche. Bitte entschuldigt, ich verspreche, dass im nĂ€chsten Teil Action und Grafiken kommen werden.
Also, im nÀchsten Teil werden wir:
- Ein kleines Client-Projekt fĂŒr alphavantage mit aiohttp und Anfragen an die benötigten Endpunkte schreiben.
- Einen Agenten erstellen, der Daten ĂŒber Wertpapiere und deren historische Preise sammelt.
Quelle: habr.com
