
¿Cómo llegué a esta vida?
No hace mucho, tuve que trabajar en el backend de un proyecto de alta carga, donde era necesario organizar la ejecución regular de una gran cantidad de tareas en segundo plano con cálculos complejos y solicitudes a servicios externos. El proyecto es asíncrono y, antes de mi llegada, tenía un mecanismo simple de cron para la ejecución de tareas: un ciclo que verificaba la hora actual y lanzaba grupos de corrutinas a través de gather; este enfoque fue aceptable hasta que las corrutinas sumaron decenas y cientos, sin embargo, cuando su número superó las dos mil, hubo que pensar en organizar una cola de tareas adecuada con un intermediario, varios trabajadores y demás.
Primero decidí probar Celery, del que había usado anteriormente. Dada la asincronía del proyecto, profundicé en el tema y vi , así como , creado por el autor del artículo.
Diré que el proyecto es muy interesante y funciona con éxito en otras aplicaciones de nuestro equipo, y el propio autor menciona que pudo implementarlo en producción, utilizando un pool asíncrono. Pero, desafortunadamente, no se adaptó muy bien a mis necesidades, ya que se descubrió con la ejecución grupal de tareas (ver ). Al momento de redactar este artículo ya estaba cerrado, sin embargo, el trabajo se llevó a cabo durante un mes. En cualquier caso, le deseo suerte al autor y lo mejor, ya que hay funcionalidades que funcionan en la biblioteca... en resumen, el problema soy yo y para mí la herramienta resultó un poco cruda. Además, en algunas tareas había de 2 a 3 solicitudes http a diferentes servicios, así que incluso optimizando las tareas, generamos 4 mil conexiones TCP aproximadamente cada 2 horas; no es lo ideal... Me gustaría crear una sesión para un tipo de tarea al iniciar los trabajadores. Un poco más sobre la gran cantidad de solicitudes a través de aiohttp .
Con esto en mente, comencé a buscar alternativas ¡y encontré! Los creadores de Celery, específicamente, como entendí , crearon , originalmente para el proyecto . Faust está inspirada en Kafka Streams y trabaja con Kafka como intermediario, además de usar rocksdb para almacenar los resultados del trabajo de los agentes, y lo más importante, es que la biblioteca es asíncrona.
También puedes ver celery y faust de los creadores de la última: sus diferencias, diferencias de los brokers, implementación de una tarea elemental. Todo es bastante simple, sin embargo, en faust destaca una característica agradable: datos tipados para la transmisión en el tópico.
¿Qué vamos a hacer?
Así que, en esta pequeña serie de artículos, mostraré cómo recoger datos en tareas en segundo plano utilizando Faust. La fuente para nuestro proyecto de ejemplo será, como se deduce del título, . Demostraré cómo escribir agentes (sink, tópicos, particiones), cómo hacer ejecuciones regulares (cron), los comandos cli más cómodos de faust (una envoltura sobre click), un clustering simple, y al final integraremos datadog (funcionando de forma predeterminada) y trataremos de ver algo. Para almacenar los datos recopilados utilizaremos mongodb y motor para la conexión.
P.D. Dada la confianza con la que se ha escrito el apartado sobre monitoreo, creo que al final del último artículo el lector se verá algo así:

Requisitos del proyecto
Dado que ya he prometido algunas cosas, hagamos una pequeña lista de lo que debe ser capaz de hacer el servicio:
- Descargar valores mobiliarios y un resumen sobre ellos (incluidos ganancias y pérdidas, balance, flujo de efectivo — del último año) — regularmente
- Descargar datos históricos (encontrar los extremos del precio de cierre por cada año comercial) — regularmente
- Descargar los últimos datos comerciales — regularmente
- Descargar una lista configurada de indicadores para cada valor mobiliario — regularmente
Como corresponde, elegimos un nombre para el proyecto al azar: horton
Preparar la infraestructura
El encabezado es fuerte, sin embargo, todo lo que necesitamos hacer es escribir un pequeño config para docker-compose con kafka (y zookeeper — en un solo contenedor), kafdrop (si queremos ver los mensajes en los tópicos), mongodb. Obtenemos [docker-compose.yml]() del siguiente tipo:
versión: '3'
servicios:
db:
container_name: horton-mongodb-local
image: mongo:4.2-bionic
command: mongod --port 20017
restart: siempre
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: siempre
ports:
- "2181:2181"
- "9092:9092"
environment:
KAFKA_LISTENERS: "INTERNO://:29092,EXTERNO://:9092"
KAFKA_ADVERTISED_LISTENERS: "INTERNO://kafka-service:29092,EXTERNO://localhost:9092"
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: "INTERNO:PLAINTEXT,EXTERNO:PLAINTEXT"
KAFKA_INTER_BROKER_LISTENER_NAME: "INTERNO"
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: siempre
ports:
- '9000:9000'
environment:
KAFKA_BROKERCONNECT: kafka-service:29092
depends_on:
- kafka-serviceAquí no hay nada complicado. Para kafka se declararon dos listeners: uno (interno) para uso dentro de la red compuesta, y el segundo (externo) para solicitudes externas, por lo que lo abrimos al exterior. 2181 es el puerto del zookeeper. Por lo demás, creo que está claro.
Preparando la estructura del proyecto
En la versión básica, la estructura de nuestro proyecto debería verse así:
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 **Todo lo que he señalado lo dejamos por ahora, y simplemente creamos archivos vacíos.**
Hemos creado la estructura. Ahora agreguemos las dependencias necesarias, redactaremos la configuración y la conexión a mongodb. No incluiré el texto completo de los archivos en el artículo para no extenderme, sino que haré enlaces a las versiones necesarias.
Comencemos con las dependencias y los metadatos del proyecto —
A continuación, iniciamos la instalación de las dependencias y la creación de virtualenv (o pueden crear manualmente la carpeta venv y activar el entorno):
pip3 install poetry (si aún no está instalado)
poetry installAhora crearemos — credenciales y a dónde conectarnos. De una vez podemos colocar los datos para alphavantage allí. Y continuamos con — extraemos datos para la aplicación de nuestra configuración. Sí, confieso, utilicé mi librería — .
La conexión con mongo es muy sencilla. Declaramos para la conexión y para los CRUD, para facilitar las solicitudes a las colecciones.
¿Qué seguirá después?
El artículo no es muy extenso, ya que aquí solo hablo sobre motivación y preparación, así que no se lo tomen a mal; prometo que en la siguiente parte habrá acción y gráficos.
Así que, en esta misma próxima parte vamos a:
- Escribir un pequeño cliente para alphavantage en aiohttp con solicitudes a los puntos finales que nos interesan.
- Crear un agente que recogerá datos sobre valores y sus precios históricos.
Fuente: habr.com
