Фонови задачи на Фауст, Част I: Въведение

Фонови задачи на Фауст, Част I: Въведение

Как стигнах до такъв живот?

Неотдавна ми се наложи да работя върху бекенда на високо натоварен проект, в който трябваше да организирам редовно изпълнение на голямо количество фонови задачи с сложни изчисления и запитвания към външни услуги. Проектът е асинхронен и преди да дойда, в него имаше прост механизъм за стартиране на задачи с cron: цикъл с проверка на текущото време и стартиране на групи корутини чрез gather — този подход бе приемлив до момента, в който броят на корутините стана десетки и стотици, но когато надхвърли две хиляди, се наложи да се помисли за организиране на нормална опашка за задачи с брокер, няколко работника и така нататък.

Първо реших да изпробвам Celery, който бях използвал преди. Поради асинхронността на проекта, се потопих в темата и видях статия, а също така проект, създаден от автора на статията.

Ще кажа така, проектът е много интересен и работи със сполука в други приложения на нашия екип, а и самият автор говори за това, че е успял да пусне в продакшън, използвайки асинхронен пул. Но, за съжаление, това не ми подхожда много, тъй като се появи проблем проблем със груповото стартиране на задачи (виж. група). В момента на написване на статията issue вече е затворена, но работата беше извършвана в продължение на месец. Във всеки случай, пожелавам успех на автора и всичко най-добро, тъй като работещи неща на библиотеката вече има… всъщност, става дума за мен и инструментът ми се стори малко недозряло. Освен това, в някои задачи имаше по 2-3 http запитвания към различни услуги, така че дори при оптимизация на задачите създаваме 4000 tcp връзки на всеки 2 часа — не е много удобно… Бих искал да създавам сесия за един тип задачи при стартиране на работниците. По-подробно за голямото количество запитвания чрез aiohttp тук.

Във връзка с това, започнах да търся алтернативи и намерих! Създаването на celery, а конкретно, както разбрах Ask Solem, беше създадено Faust, първоначално за проекта robinhood. Faust е написана под впечатление от Kafka Streams и работи с Kafka като брокер, а за съхранение на резултатите от работата на агентите се използва rocksdb, а най-важното — библиотеката е асинхронна.

Също така, можете да разгледате кратко сравнение celery и faust от създателите на последната: разликите, различията на брокерите, реализирането на елементарна задача. Всичко е доста просто, обаче, faust привлекателно допълва с приятна характеристика — типизирани данни за предаване в топик.

Какво ще правим?

И така, в малка серия статии ще покажа как да събираш данни в фонови задачи с помощта на Faust. Източникът на нашия примерен проект ще бъде, както подсказва заглавието, alphavantage.co. Ще демонстрирам как да пишем агенти (sink, топици, партиции), как да правим редовно (cron) изпълнение, удобни cli команди на faust (обвивка над click), прост клъстеринг и в края ще свържем datadog (работещ от кутията) и ще се опитаме да видим нещо. За съхранение на събраните данни ще използваме mongodb и motor за свързване.

P.S. Съдейки по увереността, с която е написан пунктът за мониторинг, мисля, че читателят в края на последната статия все пак ще изглежда, по някакъв начин така:

Фонови задачи на Фауст, Част I: Въведение

Изисквания към проекта

С оглед на това, че вече съм успял да обещая, ще съставим малък списък с това, което услугата трябва да може:

  1. Да изтегля ценни книжа и обобщение за тях (в т.ч. печалби и загуби, баланс, cash flow — за последната година) — редовно
  2. Да изтегля исторически данни (за всяка търговска година да намери екстремумите на цената на затваряне на търговията) — редовно
  3. Да изтегля последни търговски данни — редовно
  4. Да изтегля конфигуриран списък с индикатори за всяка ценна книга — редовно

Както се полага, избираме името на проекта от нищото: horton

Подготвяме инфраструктурата

Заглавието, разбира се, е силно, но всичко, което трябва да направим, е да напишем малък конфиг за docker-compose с kafka (и zookeeper — в един контейнер), kafdrop (ако искаме да видим съобщенията в топиците), mongodb. Получаваме [docker-compose.yml](https://github.com/Egnod/horton/blob/562fa5ec14df952cd74760acf76e141707d2ef58/docker-compose.yml) следния вид:

версия: '3'

услуги:
  db:
    имя_контейнера: horton-mongodb-local
    изображение: mongo:4.2-bionic
    команда: mongod --port 20017
    перезапуск: всегда
    порты:
      - 20017:20017
    окружение:
      - MONGO_INITDB_DATABASE=horton
      - MONGO_INITDB_ROOT_USERNAME=admin
      - MONGO_INITDB_ROOT_PASSWORD=admin_password

  kafka-service:
    имя_контейнера: horton-kafka-local
    изображение: obsidiandynamics/kafka
    перезапуск: всегда
    порты:
      - "2181:2181"
      - "9092:9092"
    окружение:
      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:
    имя_контейнера: horton-kafdrop-local
    изображение: 'obsidiandynamics/kafdrop:latest'
    перезапуск: всегда
    порты:
      - '9000:9000'
    окружение:
      KAFKA_BROKERCONNECT: kafka-service:29092
    зависит_от:
      - kafka-service

Тук наистина няма нищо сложно. За kafka обявихме два listener-а: един (internal) за използване вътре в композната мрежа, а втория (external) за заявки от вън, затова го проброшихме навън. 2181 е портът на zookeeper-а. По останалото, мисля, е ясно.

Готвим скелет на проекта

В основния вариант структурата на нашия проект трябва да изглежда така:

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 *

*Всичко, което отбелязах все още не докосваме, а просто създаваме празни файлове.**

Създадохме структура. Сега да добавим необходимите зависимости, да напишем конфиг и да свържем с mongodb. Пълен текст на файловете няма да предоставя в статията, за да не забавям, а ще дам линкове към необходимите версии.

Ще започнем с зависимостите и мета информацията за проекта — pyproject.toml

След това стартираме инсталацията на зависимостите и създаването на virtualenv (или, можете сами да създадете папка venv и да активирате околната среда):

pip3 install poetry (ако все още не е инсталирано)
poetry install

Сега да създадем config.yml — креденциали и къде да се свързваме. Веднага можем да сложим и данните за alphavantage. И преминаваме към config.py — извличаме данни за приложението от нашия конфиг. Да, признавам, използвах своята библиотека — sitri.

По свързването с mongodb — всичко е наистина просто. Обявихме клас клиент за свързване и основен клас за CRUD операции, за по-лесно правене на заявки по колекциите.

Какво ще последва?

Статията не е много дълга, тъй като тук говоря само за мотивацията и подготовката, така че не се сърдете — обещавам, че в следващата част ще има действие и графики.

И така, в тази следваща част ние:

  1. Ще напишем малък клиент за alphavantage на aiohttp с заявки към необходимите ни ендпойнти.
  2. Ще създадем агент, който ще събира данни за ценни книжа и исторически цени за тях.

Код на проекта

Кодът на тази част

Източник: habr.com

Купете надежден хостинг за сайтове със защита от DDoS, VPS и VDS сървъри 🔥 Купете надежден хостинг за сайтове със защита от DDoS, VPS и VDS сървъри | ProHoster