Achtergrondtaken op Faust, Deel I: Inleiding

Achtergrondtaken op Faust, Deel I: Inleiding

Hoe ben ik hier beland?

Niet zo lang geleden moest ik werken aan de backend van een zwaarbelast project, waarbij een regelmatige uitvoering van een groot aantal achtergrondtaken met complexe berekeningen en verzoeken naar externe services moest worden georganiseerd. Het project is asynchroon en voordat ik kwam, was er een eenvoudige cron-methode voor het starten van taken: een loop met een controle van de huidige tijd en het starten van groepen coroutines via gather — deze aanpak bleek acceptabel totdat het aantal coroutines in de tientallen en honderden begon te lopen, maar toen het aantal meer dan tweeduizend overschreed, moest ik gaan nadenken over een goede organisatie van een taakqueue met een broker, meerdere workers en andere zaken.

Eerst besloot ik Celery uit te proberen, wat ik eerder had gebruikt. Gezien de asynchroniciteit van het project, dook ik in het onderwerp en zag artikel, evenals project, gemaakt door de auteur van het artikel.

Ik zal het zo zeggen, het project is erg interessant en werkt behoorlijk succesvol in andere toepassingen van ons team, en de auteur zelf zegt ook dat hij het in productie heeft kunnen uitrollen, gebruikmakend van een asynchrone pool. Maar helaas paste het niet zo goed bij me, omdat er een probleem was met het groepsgewijs starten van taken (zie groep). Op het moment dat dit artikel werd geschreven, probleem was het al opgelost, echter, het werk liep een maand door. Hoe dan ook, veel succes en het allerbeste voor de auteur, omdat er al werkende dingen in de lib zijn... in ieder geval, het ligt aan mij en voor mij bleek de tool nog steeds niet helemaal klaar. Bovendien waren er in sommige taken 2-3 http-verzoeken naar verschillende services, hierdoor creëren we zelfs bij optimalisatie van taken 4000 tcp-verbindingen, ongeveer elke 2 uur — dat is niet echt ideaal... Het zou fijn zijn om een sessie voor één type taak op te zetten bij het starten van workers. Wat betreft het grote aantal verzoeken via aiohttp, here.

daarom begon ik te zoeken naar alternatieven en ik vond het! De makers van Celery, en specifiek, zoals ik het begreep, Ask Solem, hebben Faust, oorspronkelijk gemaakt voor het project robinhood. Faust is geschreven ter inspiratie van Kafka Streams en werkt met Kafka als broker, ook wordt rocksdb gebruikt voor de opslag van resultaten van de agenten, en het belangrijkste is dat de bibliotheek asynchroon is.

Daarnaast kun je ook kijken naar een korte vergelijking Celery en Faust van de makers van de nieuwste: hun verschillen, de verschillen tussen de brokers, de implementatie van een eenvoudige taak. Het is allemaal vrij simpel, maar in Faust valt een aangename functie op — getypeerde gegevens voor overdracht naar een topic.

Wat gaan we doen?

Dus, in een kleine serie artikelen laat ik zien hoe je gegevens kunt verzamelen in achtergrondtaken met behulp van Faust. De bron voor ons voorbeeldproject zal, zoals de titel al aangeeft, alphavantage.cozijn. Ik zal demonstreren hoe je agenten (sink, topics, partitions) schrijft, hoe je regelmatig (cron) uitvoert, handige cli-commando's van Faust (een wrapper boven click), eenvoudige clustering, en aan het eind zullen we Datadog (direct werkend) aansluiten en proberen iets te zien. Voor het opslaan van de verzamelde gegevens zullen we MongoDB en Motor voor de connectie gebruiken.

P.S. Gezien het vertrouwen waarmee de sectie over monitoring is geschreven, denk ik dat de lezer aan het eind van het laatste artikel er ongeveer zo zal uitzien:

Achtergrondtaken op Faust, Deel I: Inleiding

Vereisten voor het project

Aangezien ik al wat beloofd heb, laten we een korte lijst opstellen van wat de service moet kunnen:

  1. Waarde effecten en overzicht daarvan exporteren (inclusief winst en verlies, balans, cashflow — voor het afgelopen jaar) — regelmatig
  2. Historische gegevens exporteren (extremen van de sluitprijs voor elk handelsjaar vinden) — regelmatig
  3. Laatste handelsgegevens exporteren — regelmatig
  4. Een geconfigureerde lijst van indicatoren voor elk effect exporteren — regelmatig

Zoals het hoort, kiezen we een naam voor het project uit de lucht: horton

De infrastructuur voorbereiden

De titel is natuurlijk sterk, maar alles wat we moeten doen is een kleine configuratie voor docker-compose schrijven met Kafka (en Zookeeper — in één container), Kafdrop (als we de berichten in de topics willen bekijken), MongoDB. We krijgen [docker-compose.yml](https://github.com/Egnod/horton/blob/562fa5ec14df952cd74760acf76e141707d2ef58/docker-compose.yml) van de volgende aard:

versie: '3'

diensten:
  db:
    container_name: horton-mongodb-local
    image: mongo:4.2-bionic
    command: mongod --port 20017
    restart: altijd
    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: altijd
    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: altijd
    ports:
      - '9000:9000'
    environment:
      KAFKA_BROKERCONNECT: kafka-service:29092
    depends_on:
      - kafka-service

Dit is helemaal niet ingewikkeld. Voor kafka hebben we twee listeners gedefinieerd: één (intern) voor gebruik binnen het compositienetwerk en de andere (extern) voor verzoeken van buitenaf, daarom hebben we deze naar buiten doorgestuurd. 2181 is de poort van de zookeeper. Verder denk ik dat het duidelijk is.

Laten we de basis van het project creëren.

In de basisstructuur zou ons project er als volgt uit moeten zien:

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 wat ik heb gemarkeerd laten we voorlopig nog even met rust, terwijl we gewoon lege bestanden aanmaken.**

We hebben de structuur aangemaakt. Laten we nu de benodigde afhankelijkheden toevoegen, de configuratie schrijven en de verbinding met mongodb maken. Ik zal de volledige tekst van de bestanden niet in het artikel opnemen, om het niet te langdradig te maken, maar ik zal links maken naar de benodigde versies.

Laten we beginnen met de afhankelijkheden en meta-informatie over het project — pyproject.toml

Daarna starten we de installatie van de afhankelijkheden en maken we de virtualenv aan (of u kunt zelf een map venv aanmaken en de omgeving activeren):

pip3 install poetry (als het nog niet is geïnstalleerd)
poetry install

Laten we nu config.yml — de credentials en waar te verbinden. Hier kunnen we ook de gegevens voor alphavantage plaatsen. En dan gaan we verder met config.py — we halen gegevens voor de applicatie uit onze configuratie. Ja, ik geef het toe, ik heb mijn eigen bibliotheek gebruikt — sitri.

Wat betreft de verbinding met mongo — het is helemaal niet ingewikkeld. We hebben de client klasse voor verbinding gedefinieerd en de basisklasse voor de CRUD operaties, zodat het gemakkelijker is om verzoeken naar de collecties te doen.

Wat zal er daarna gebeuren?

Het artikel is niet zo lang geworden, omdat ik hier alleen over motivatie en voorbereiding spreek, dus mijn excuses — ik beloof dat er in het volgende deel actie en graphics zullen komen.

Dus, in dit volgende deel gaan we:

  1. Een kleine client voor alphavantage schrijven met aiohttp die verzoeken naar de benodigde eindpunten verzendt.
  2. Een agent maken die gegevens over effecten verzamelt en historische prijzen daarvoor opzoekt.

De code van het project

De code van dit deel

Bron: habr.com

Koop betrouwbare webhosting met bescherming tegen DDoS, VPS VDS servers 🔥 Koop betrouwbare webhosting met bescherming tegen DDoS, VPS VDS servers | ProHoster