Compiti in background su Faust, Parte I: Introduzione

Compiti in background su Faust, Parte I: Introduzione

Come sono arrivato a questo punto?

Recentemente ho lavorato su un backend di un progetto ad alta intensità di carico, in cui era necessario organizzare l'esecuzione regolare di un gran numero di compiti in background con calcoli complessi e richieste a servizi esterni. Il progetto è asincrono e, prima del mio arrivo, utilizzava un semplice meccanismo di attivazione delle attività tramite cron: un ciclo che controllava l'ora corrente e avviava gruppi di coroutine tramite gather. Questo approccio era accettabile fino a quando il numero di coroutine era limitato a decine e centinaia; tuttavia, quando il numero ha superato le duemila, è stato necessario pensare a una normale organizzazione della coda di compiti con un broker, più worker e altro.

Inizialmente ho deciso di provare Celery, di cui avevo fatto uso in precedenza. Data l'asincronicità del progetto, mi sono immerso nella questione e ho visto articolo, e anche progetto, creato dall'autore dell'articolo.

Posso dire che il progetto è molto interessante e sta funzionando piuttosto bene in altre applicazioni del nostro team. Inoltre, l'autore stesso afferma di essere riuscito a rilasciarlo in produzione utilizzando un pool asincrono. Tuttavia, purtroppo, non si adatta molto bene a me, dato che è emerso... problema con il lancio di gruppi di attività (vedi... gruppo). Al momento della scrittura di questo articolo... issue è già chiuso, tuttavia, il lavoro è stato svolto per un mese. In ogni caso, auguro buona fortuna all'autore e ogni bene, poiché ci sono già funzionalità operative nella libreria... insomma, il problema è mio e per me lo strumento è risultato un po' grezzo. Inoltre, in alcuni compiti c'erano 2-3 richieste http a diversi servizi, quindi anche ottimizzando le attività creiamo 4.000 connessioni tcp, circa ogni 2 ore — non proprio una buona soluzione... Vorrei poter creare una sessione per un tipo di attività al momento dell'avvio dei lavoratori. Un po' più in dettaglio riguardo all'alto numero di richieste tramite aiohttp. qui.

A causa di ciò, ho iniziato a cercare... alternative e l'ho trovata! Creati dai fondatori di celery, e in particolare, come ho capito... Ask Solem, è stato creato Faust, originariamente per il progetto robinhood. Faust è ispirato a Kafka Streams e utilizza Kafka come broker; inoltre, per memorizzare i risultati delle operazioni degli agenti viene usato rocksdb. Ma la cosa più importante è che la libreria è asincrona.

Inoltre, puoi dare un'occhiata a un breve confronto tra celery e faust, da parte dei creatori dell'ultimo: le loro differenze, quelle dei broker, la realizzazione di un compito elementare. È tutto piuttosto semplice, però in faust colpisce un aspetto interessante: i dati tipizzati per il passaggio al topic.

Cosa faremo?

Quindi, in una breve serie di articoli, mostrerò come raccogliere dati in attività in background utilizzando Faust. La fonte per il nostro progetto esemplare sarà, come suggerisce il nome, alphavantage.co. Dimostrerò come scrivere agenti (sink, topic, partizioni), come effettuare esecuzioni regolari (cron), le comodissime comandi cli di faust (un wrapper su click), un semplice clustering, e infine integreremo datadog (pronto all'uso) e cercheremo di vedere qualcosa. Per la memorizzazione dei dati raccolti utilizzeremo mongodb e motor per la connessione.

P.S. Dato il livello di sicurezza con cui è scritto il punto sul monitoraggio, penso che il lettore alla fine dell'ultimo articolo si presenterà in questo modo:

Compiti in background su Faust, Parte I: Introduzione

Requisiti per il progetto

Considerando che ho già fatto alcune promesse, prepariamo una breve lista di ciò che deve fare il servizio:

  1. Scaricare titoli e panoramiche su di essi (inclusi profitti e perdite, bilancio, flussi di cassa - nell'ultimo anno) - regolarmente
  2. Scaricare dati storici (trovare i massimi e minimi del prezzo di chiusura per ogni anno di trading) - regolarmente
  3. Scaricare gli ultimi dati di trading - regolarmente
  4. Scaricare una lista configurata di indicatori per ciascun titolo - regolarmente

Come da prassi, scegliamo un nome per il progetto a caso: horton

Prepariamo l'infrastruttura

Il titolo è sicuramente forte, tuttavia, tutto ciò che dobbiamo fare è scrivere una piccola configurazione per docker-compose con kafka (e zookeeper - in un unico contenitore), kafdrop (se vogliamo dare un'occhiata ai messaggi nei topic), mongodb. Otteniamo [docker-compose.yml](https://github.com/Egnod/horton/blob/562fa5ec14df952cd74760acf76e141707d2ef58/docker-compose.yml) di questo tipo:

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-service

Non c'è davvero nulla di complesso. Per Kafka sono stati dichiarati due listener: uno (interno) per l'uso all'interno della rete composita e l'altro (esterno) per le richieste dall'esterno, motivo per cui è stato esposto. 2181 è la porta di Zookeeper. Per il resto, credo sia chiaro.

Prepariamo la struttura del progetto

Nel caso base, la struttura del nostro progetto dovrebbe apparire così:

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 *

*Tutto ciò che ho contrassegnato per ora non lo tocchiamo, ma creiamo semplicemente file vuoti.**

Abbiamo creato la struttura. Ora aggiungiamo le dipendenze necessarie, scriviamo la configurazione e ci connettiamo a mongodb. Non riporterò il testo completo dei file nell'articolo per non allungarlo, ma metterò i collegamenti alle versioni necessarie.

Iniziamo con le dipendenze e le informazioni sul progetto — pyproject.toml

Dopo, avviamo l'installazione delle dipendenze e la creazione del virtualenv (oppure, potete creare la cartella venv e attivare l'ambiente):

pip3 install poetry (se non è già installato)
poetry install

Ora creiamo config.yml — credenziali e a chi connettersi. Possiamo subito inserire anche i dati per alphavantage. Ora passiamo a config.py — estraiamo i dati per l'applicazione dalla nostra configurazione. Sì, confesso, ho usato la mia libreria — sitri.

Per connettersi a mongo — è davvero semplice. Abbiamo dichiarato una classe cliente per la connessione e una classe base per i crud, per facilitare le richieste sulle collezioni.

Cosa succederà dopo?

L'articolo non è molto lungo, in quanto parlo solo di motivazione e preparazione, quindi non prendetevela — prometto che nella prossima parte ci sarà azione e grafica.

Quindi, in questa prossima parte faremo:

  1. Scriveremo un piccolo client per alphavantage usando aiohttp con richieste ai punti finali di cui abbiamo bisogno.
  2. Creeremo un agente che raccoglierà dati sui titoli e i loro prezzi storici.

Codice del progetto

Il codice di questa parte

Fonte: habr.com

Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server 🔥 Acquista hosting affidabile per siti web con protezione DDoS, VPS VDS server | ProHoster