Attività in background su Faust, Parte I: Introduzione

Attività in background su Faust, Parte I: Introduzione

Come sono arrivato a una vita del genere?

Non molto tempo fa mi è capitato di lavorare sul backend di un progetto ad alta intensità di carico, nel quale era necessario organizzare l'esecuzione regolare di un gran numero di attività in background con calcoli complessi e richieste a servizi esterni. Il progetto è asincrono e, prima del mio arrivo, c'era un semplice meccanismo di cron per l'avvio delle attività: un ciclo con il controllo del tempo attuale e l'avvio di gruppi di coroutine tramite gather — questo approccio si era dimostrato accettabile fino a quando il numero di coroutine era di alcune decine o centinaia, tuttavia, quando il loro numero ha superato le duemila, è stato necessario pensare a un'organizzazione di una normale coda di attività con un broker, diversi 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, così come progetto, creato dall'autore dell'articolo.

Dico solo che il progetto è molto interessante e funziona abbastanza bene in altre applicazioni del nostro team, e l'autore stesso dice di essere riuscito a lanciarlo in produzione, utilizzando un pool asincrono. Ma, sfortunatamente, non mi è sembrato adatto, poiché sono emerse il problema problematiche con l'avvio di attività di gruppo (vedi gruppo). Al momento della scrittura dell'articolo problema era già chiuso, tuttavia, il lavoro è stato svolto per un mese. In ogni caso, auguro all'autore buona fortuna e ogni bene, poiché ci sono già cose funzionanti nella libreria… insomma, è un problema mio e si è rivelato uno strumento un po' grezzo per me. Inoltre, in alcune attività c'erano da 2 a 3 richieste HTTP a diversi servizi, creando così circa 4000 connessioni TCP ogni 2 ore — non è proprio ideale… Mi piacerebbe creare una sessione per un tipo di attività all'avvio dei worker. Maggiori dettagli su un gran numero di richieste tramite aiohttp qui.

A questo proposito, ho iniziato a cercare alternative e ho trovato! Dagli autori di Celery, e in particolare, come ho capito Ask Solem, è stato creato Faust, inizialmente per il progetto robinhood. Faust è scritto sotto l'influenza di Kafka Streams e funziona con Kafka come broker, inoltre, per la memorizzazione dei risultati del lavoro degli agenti, viene utilizzato rocksdb, e la cosa più importante è che la libreria è asincrona.

Puoi anche dare un'occhiata a un breve confronto celery e faust dei creatori dell'ultimo: le loro differenze, le differenze tra broker, l'implementazione di un compito elementare. Tutto è piuttosto semplice, tuttavia, in faust attira l'attenzione una piacevole caratteristica: dati tipizzati per la trasmissione nel topic.

Cosa facciamo?

Quindi, in una breve serie di articoli, mostrerò come raccogliere dati in compiti in background tramite Faust. La fonte del nostro progetto esemplificativo sarà, come suggerisce il nome, alphavantage.co. Dimostrerò come scrivere agenti (sink, topic, partizioni), come fare esecuzioni regolari (cron), comandi cli molto convenienti di faust (un wrapper su click), clustering semplice, e alla fine attaccheremo datadog (che funziona out of the box) e cercheremo di vedere qualcosa. Per memorizzare i dati raccolti useremo mongodb e motor per la connessione.

P.S. A giudicare dalla sicurezza con cui è scritto il punto sul monitoraggio, penso che il lettore alla fine dell'ultimo articolo comunque apparirà qualcosa di simile a questo:

Attività in background su Faust, Parte I: Introduzione

Requisiti del progetto

Dato che ho già promesso, compiliamo un breve elenco di ciò che il servizio deve saper fare:

  1. Scaricare titoli e overview su di essi (inclusi profitti e perdite, bilancio, cash flow - per l'ultimo anno) - regolarmente
  2. Scaricare dati storici (per ogni anno di trading trovare i picchi del prezzo di chiusura del trading) - regolarmente
  3. Scaricare gli ultimi dati di trading - regolarmente
  4. Scaricare un elenco configurato di indicatori per ogni titolo - regolarmente

Come si conviene, scegliamo il nome del progetto a caso: horton

Prepariamo l'infrastruttura

Il titolo è ovviamente forte, tuttavia, tutto ciò che deve essere fatto è scrivere un piccolo file di configurazione per docker-compose con kafka (e zookeeper - in un solo contenitore), kafdrop (se vogliamo vedere i messaggi nei topic), mongodb. Otteniamo [docker-compose.yml](https://github.com/Egnod/horton/blob/562fa5ec14df952cd74760acf76e141707d2ef58/docker-compose.yml) della seguente forma:

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'è nulla di complicato. 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, quindi è stato esposto. 2181 è la porta di Zookeeper. Per il resto, credo sia chiaro.

Prepariamo la struttura del progetto

Nella versione di 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 evidenziato non lo tocchiamo per ora, ma creiamo semplicemente file vuoti.**

Abbiamo creato la struttura. Ora aggiungiamo le dipendenze necessarie, scriviamo la configurazione e ci colleghiamo a MongoDB. Non riporterò il testo completo dei file nell'articolo per non allungare, ma farò dei link alle versioni necessarie.

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

Poi, avviamo l'installazione delle dipendenze e la creazione dell'ambiente virtuale (oppure, potete creare una cartella venv e attivare l'ambiente da soli):

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

Ora creiamo config.yml — le credenziali e dove connettersi. Possiamo già inserire anche i dati per Alphavantage. Adesso passiamo a config.py — estraiamo i dati per l'applicazione dalla nostra configurazione. Sì, lo ammetto, ho utilizzato la mia libreria — sitri.

Per la connessione con MongoDB — è tutto molto semplice. Abbiamo dichiarato la classe client per la connessione e la classe base per i CRUD, per facilitare le richieste sulle collezioni.

Cosa succederà dopo?

L'articolo è piuttosto breve, poiché qui parlo solo di motivazione e preparazione, quindi non me ne vogliate — prometto che nella prossima parte ci sarà azione e grafica.

Quindi, in questa prossima parte faremo:

  1. Scriveremo un piccolo client per alphavantage utilizzando aiohttp con richieste agli endpoint 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