
Как стигнах до такъв живот?
Неотдавна ми се наложи да работя върху бекенда на високо натоварен проект, в който трябваше да организирам редовно изпълнение на голямо количество фонови задачи с сложни изчисления и запитвания към външни услуги. Проектът е асинхронен и преди да дойда, в него имаше прост механизъм за стартиране на задачи с cron: цикъл с проверка на текущото време и стартиране на групи корутини чрез gather — този подход бе приемлив до момента, в който броят на корутините стана десетки и стотици, но когато надхвърли две хиляди, се наложи да се помисли за организиране на нормална опашка за задачи с брокер, няколко работника и така нататък.
Първо реших да изпробвам Celery, който бях използвал преди. Поради асинхронността на проекта, се потопих в темата и видях , а също така , създаден от автора на статията.
Ще кажа така, проектът е много интересен и работи със сполука в други приложения на нашия екип, а и самият автор говори за това, че е успял да пусне в продакшън, използвайки асинхронен пул. Но, за съжаление, това не ми подхожда много, тъй като се появи проблем със груповото стартиране на задачи (виж. ). В момента на написване на статията вече е затворена, но работата беше извършвана в продължение на месец. Във всеки случай, пожелавам успех на автора и всичко най-добро, тъй като работещи неща на библиотеката вече има… всъщност, става дума за мен и инструментът ми се стори малко недозряло. Освен това, в някои задачи имаше по 2-3 http запитвания към различни услуги, така че дори при оптимизация на задачите създаваме 4000 tcp връзки на всеки 2 часа — не е много удобно… Бих искал да създавам сесия за един тип задачи при стартиране на работниците. По-подробно за голямото количество запитвания чрез aiohttp .
Във връзка с това, започнах да търся алтернативи и намерих! Създаването на celery, а конкретно, както разбрах , беше създадено , първоначално за проекта . Faust е написана под впечатление от Kafka Streams и работи с Kafka като брокер, а за съхранение на резултатите от работата на агентите се използва rocksdb, а най-важното — библиотеката е асинхронна.
Също така, можете да разгледате celery и faust от създателите на последната: разликите, различията на брокерите, реализирането на елементарна задача. Всичко е доста просто, обаче, faust привлекателно допълва с приятна характеристика — типизирани данни за предаване в топик.
Какво ще правим?
И така, в малка серия статии ще покажа как да събираш данни в фонови задачи с помощта на Faust. Източникът на нашия примерен проект ще бъде, както подсказва заглавието, . Ще демонстрирам как да пишем агенти (sink, топици, партиции), как да правим редовно (cron) изпълнение, удобни cli команди на faust (обвивка над click), прост клъстеринг и в края ще свържем datadog (работещ от кутията) и ще се опитаме да видим нещо. За съхранение на събраните данни ще използваме mongodb и motor за свързване.
P.S. Съдейки по увереността, с която е написан пунктът за мониторинг, мисля, че читателят в края на последната статия все пак ще изглежда, по някакъв начин така:

Изисквания към проекта
С оглед на това, че вече съм успял да обещая, ще съставим малък списък с това, което услугата трябва да може:
- Да изтегля ценни книжа и обобщение за тях (в т.ч. печалби и загуби, баланс, cash flow — за последната година) — редовно
- Да изтегля исторически данни (за всяка търговска година да намери екстремумите на цената на затваряне на търговията) — редовно
- Да изтегля последни търговски данни — редовно
- Да изтегля конфигуриран списък с индикатори за всяка ценна книга — редовно
Както се полага, избираме името на проекта от нищото: horton
Подготвяме инфраструктурата
Заглавието, разбира се, е силно, но всичко, което трябва да направим, е да напишем малък конфиг за docker-compose с kafka (и zookeeper — в един контейнер), kafdrop (ако искаме да видим съобщенията в топиците), mongodb. Получаваме [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. Пълен текст на файловете няма да предоставя в статията, за да не забавям, а ще дам линкове към необходимите версии.
Ще започнем с зависимостите и мета информацията за проекта —
След това стартираме инсталацията на зависимостите и създаването на virtualenv (или, можете сами да създадете папка venv и да активирате околната среда):
pip3 install poetry (ако все още не е инсталирано)
poetry installСега да създадем — креденциали и къде да се свързваме. Веднага можем да сложим и данните за alphavantage. И преминаваме към — извличаме данни за приложението от нашия конфиг. Да, признавам, използвах своята библиотека — .
По свързването с mongodb — всичко е наистина просто. Обявихме за свързване и за CRUD операции, за по-лесно правене на заявки по колекциите.
Какво ще последва?
Статията не е много дълга, тъй като тук говоря само за мотивацията и подготовката, така че не се сърдете — обещавам, че в следващата част ще има действие и графики.
И така, в тази следваща част ние:
- Ще напишем малък клиент за alphavantage на aiohttp с заявки към необходимите ни ендпойнти.
- Ще създадем агент, който ще събира данни за ценни книжа и исторически цени за тях.
Източник: habr.com
