
Comment suis-je arrivé à cette vie?
Il n'y a pas si longtemps, j'ai dû travailler sur le backend d'un projet à fort trafic, où il fallait organiser l'exécution régulière d'un grand nombre de tâches de fond avec des calculs complexes et des requêtes vers des services externes. Le projet est asynchrone et avant mon arrivée, il y avait un mécanisme simple de planification des tâches par cron : une boucle vérifiant l'heure actuelle et lançant des groupes de coroutines via gather — cette approche était acceptable jusqu'à ce qu'il y ait des dizaines et des centaines de ces coroutines, mais lorsque leur nombre a dépassé deux mille, il a fallu penser à une organisation normale d'une file d'attente de tâches avec un broker, plusieurs workers et autres.
Au début, j'ai décidé d'essayer Celery, que j'avais utilisé auparavant. Étant donné l'asynchronie du projet, je me suis plongé dans la question et j'ai vu , ainsi que , créé par l'auteur de l'article.
Je dirai que le projet est très intéressant et fonctionne assez bien dans d'autres applications de notre équipe, et l'auteur lui-même dit avoir pu le déployer en production en utilisant un pool asynchrone. Mais, malheureusement, cela ne m'a pas beaucoup convenu, car il s'est avéré avec le lancement groupé des tâches (voir ). Au moment de la rédaction de l'article était déjà clôturé, mais le travail s'est poursuivi pendant un mois. Quoi qu'il en soit, je souhaite bonne chance à l'auteur et tout le meilleur, car il y a déjà des choses fonctionnelles dans la librairie... en gros, c’est de ma faute et l'outil s'est avéré un peu brut pour moi. De plus, dans certaines tâches, il y avait 2-3 requêtes HTTP vers différents services, ainsi même en optimisant les tâches, nous créons 4000 connexions TCP environ toutes les 2 heures — ce n’est pas idéal... J'aimerais établir une session pour un type de tâche lors du lancement des workers. Un peu plus de détails sur le grand nombre de requêtes via aiohttp .
En conséquence, j'ai commencé à chercher des alternatives et j'en ai trouvé une ! Créée par les développeurs de Celery, et plus précisément, si j'ai bien compris , il a été créé , à l'origine pour le projet . Faust est inspiré par Kafka Streams et fonctionne avec Kafka en tant que broker, de plus, pour le stockage des résultats du travail des agents, rocksdb est utilisé, et le plus important est que la bibliothèque est asynchrone.
Vous pouvez également consulter Celery et Faust des créateurs de la dernière : leurs différences, les différences entre les brokers, l'implémentation d'une tâche élémentaire. Tout cela est plutôt simple, cependant, Faust présente une caractéristique intéressante - des données typées à transmettre dans un topic.
Que allons-nous faire ?
Ainsi, dans une petite série d'articles, je vais montrer comment collecter des données dans des tâches en arrière-plan en utilisant Faust. La source de notre projet d'exemple sera, comme l'indique le titre, . Je vais démontrer comment écrire des agents (sink, topics, partitions), comment effectuer une exécution régulière (cron), les commandes cli de Faust très pratiques (un wrapper autour de click), un clustering simple, et à la fin, nous intégrerons Datadog (fonctionnant par défaut) et tenterons de voir quelque chose. Pour stocker les données collectées, nous utiliserons MongoDB et Motor pour la connexion.
P.S. Vu la confiance avec laquelle le point sur la surveillance est écrit, je pense que le lecteur, à la fin du dernier article, ressemblera à quelque chose comme ça :

Exigences du projet
Étant donné que j'ai déjà promis pas mal de choses, établissons une petite liste de ce que le service doit pouvoir faire :
- Exporter les titres et un aperçu de ceux-ci (y compris les gains et les pertes, le bilan, le cash flow - pour la dernière année) - régulièrement
- Exporter des données historiques (trouver les extrêmes des prix de clôture pour chaque année de trading) - régulièrement
- Exporter les dernières données de trading - régulièrement
- Exporter la liste configurée des indicateurs pour chaque titre - régulièrement
Comme il se doit, nous choisissons un nom de projet au hasard : horton
Préparons l'infrastructure
Le titre est bien sûr percutant, cependant, tout ce que nous devons faire est d'écrire une petite configuration pour docker-compose avec Kafka (et Zookeeper - dans un seul conteneur), Kafdrop (si nous souhaitons voir les messages dans les topics), MongoDB. Nous obtenons [docker-compose.yml]() de la forme suivante :
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-serviceIl n'y a rien de compliqué ici. Pour Kafka, nous avons déclaré deux écouteurs : l'un (interne) pour une utilisation au sein du réseau composite, et le second (externe) pour les requêtes externes, donc nous l'avons exposé à l'extérieur. 2181 est le port de Zookeeper. Pour le reste, je pense que c'est clair.
Préparons le squelette du projet
Dans la version de base, la structure de notre projet devrait ressembler à ceci :
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 **Tout ce que j'ai marqué nous ne le touchons pas encore, mais nous créons simplement des fichiers vides.**
La structure est créée. Maintenant, ajoutons les dépendances nécessaires, écrivons la configuration et connectons-nous à MongoDB. Je ne vais pas fournir le texte complet des fichiers dans l'article pour ne pas alourdir, mais je ferai des liens vers les versions nécessaires.
Commençons par les dépendances et les métadonnées du projet —
Ensuite, nous lançons l'installation des dépendances et la création de virtualenv (ou, vous pouvez créer le dossier venv vous-même et activer l'environnement) :
pip3 install poetry (si ce n'est pas déjà installé)
poetry installMaintenant, créons — les crédences et où se connecter. On peut également y placer les données pour Alphavantage. Et nous passons à — nous extrayons les données pour l'application depuis notre configuration. Oui, je l'avoue, j'ai utilisé ma librairie — .
Concernant la connexion avec MongoDB — c'est tout simple. Nous avons déclaré pour la connexion et pour les CRUD, afin de faciliter les requêtes sur les collections.
Que se passera-t-il ensuite ?
L'article n'est pas très long, car ici je parle seulement de la motivation et de la préparation, alors ne m'en voulez pas — je promets que dans la prochaine partie il y aura de l'action et des graphiques.
Alors, dans cette prochaine partie nous allons :
- Écrire un petit client pour alphavantage sur aiohttp avec des requêtes vers les points de terminaison dont nous avons besoin.
- Créer un agent qui collectera des données sur les valeurs mobilières et leurs prix historiques.
Source : habr.com
