Системный дизайн: сервис скрейпинг-джоб

ЗаданиеProduct Owner

Нужна система, в которой пользователи запускают скрейпинг-джобы против разных API и не ждут результата. Пользователь должен видеть, что происходит с его джобами, и получить данные, когда они готовы. API бывают медленные и ненадёжные, поэтому джобы не должны падать просто так. И пользоваться этим будут многие — система должна тянуть много джоб одновременно.

Дальше другой формат: вместо разбора одного паттерна — задача целиком. Задание выше пришло в том виде, в каком его дают на собеседовании или присылают в первом письме от продукта: абзац текста без единого числа.

Хочется сразу нарисовать схему: очередь, воркеры, база, стрелки. Схема будет, но в конце.

Очередь на схеме сама по себе ничего не говорит. Смысл появляется вместе с числами: сколько джоб в секунду она принимает, сколько их ждёт на пике, сколько джоба стоит в очереди до первого воркера. Без чисел выбор между Kafka и Postgres делается наугад.

Поэтому порядок такой:

  1. Понять, что просят. Пересказать задачу как набор обещаний системы.
  2. Спросить. Выяснить всё, чего нет в постановке, но что нужно для решения.
  3. Записать требования. Функциональные — что система умеет; нефункциональные — с какими свойствами.
  4. Посчитать. Поток, одновременность, объём данных.
  5. Нарисовать. C1, C2 и ниже, обосновывая каждое решение числами.

Каждый шаг опирается на ответы предыдущего.

  1. Шаг 1 из 36

    Что на самом деле просят

    С такой задачей сначала пересказывают её своими словами и сверяют пересказ с тем, кто её принёс. В абзаце выше нет чисел, но есть четыре обещания, и за каждое придётся платить.

    • Принять и отпустить. Пользователь отправляет джобу и сразу получает ответ: принято, вот идентификатор. Результата в этом ответе нет. Значит, приём и выполнение разнесены во времени, и между ними что-то стоит.
    • Показать, что происходит. У джобы есть состояние: ждёт, выполняется, готова, упала. Его нужно где-то хранить, кто-то должен его менять и кто-то — показывать.
    • Отдать данные, когда готовы. Результат забирают позже, отдельным запросом. Значит, он хранится вне процесса, который его получил.
    • Не падать без причины. Чужой API бывает медленным или недоступным, и это обычная ситуация. Нужны таймауты, повторы и способ отличить «не получилось сейчас» от «не получится никогда».

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

    В условии нет чисел, нет ни слова о том, что именно скрейпят и сколько хранить результат. Это нужно выяснить у продукта, и это следующий шаг.

  2. Шаг 2 из 36

    Вопросы, без ответов на которые считать нечего

    Под такую формулировку подходит и одна машина с горутинами, и сотня воркеров в трёх регионах. Какой вариант нужен, решают ответы, которых в задании нет. Поэтому дальше вопросы.

    • Что скрейпим? Публичные API с ключом и лимитом, чужие сайты без договора или выгрузки партнёров, которые нас ждут. От этого зависит главный риск: счёт за превышение, бан по IP или ни то, ни другое.
    • Кто ставит джобы? Человек в интерфейсе и чужой скрипт по-разному забирают результат.
    • Сколько джоб? «Много» может значить сотню в сутки или миллион в час. Пока число не названо, требование нельзя проверить.
    • Сколько хранить результат? Час, месяц или до ручного удаления. От этого зависят выбор хранилища и его стоимость.
    • Что считать успехом? Любой ответ цели или только пригодные данные. Во втором случае нужна валидация, и она тоже может отказать.
    • Нужны ли отмена и повтор? Для отмены воркер должен уметь прервать работу. Для повтора нужно решить, что делать с прошлым результатом.
    • Бывают ли срочные джобы? Приоритеты усложняют планировщик, и добавлять их заранее не стоит.

    Словесные ответы пойдут в требования, числовые — в расчёты. Если продукт отвечает «не знаю», записываем допущение и проверяем его числом на шаге расчётов.

  3. Шаг 3 из 36

    Кто пользуется системой

    У такой системы обычно три типа пользователей с разными потребностями.

    • Человек в интерфейсе. Запустил джобу и смотрит в список. Ему нужно видеть, как меняется состояние, и открывать готовый результат в один щелчок. Ждать он готов, если видит, что работа идёт.
    • Чужой скрипт. Ставит тысячу джоб и уходит. Опрашивать каждую по отдельности ему неудобно: нужен вебхук или запрос, который вернёт состояние многих джоб сразу.
    • Дежурный. Приходит, когда что-то сломалось. Ему нужна история конкретной джобы: на каком шаге она остановилась и почему.

    Отсюда требование к схеме: у джобы есть идентификатор, по которому все трое читают её историю, и эта история хранится в одном месте. Интерфейсы разные, состояние одно.

  4. Шаг 4 из 36

    Функциональные требования

    Требования хранятся файлом в репозитории сервиса, как словарь в главе про DDD. Вики со временем расходится с тем, что делает сервис. Файл рядом с кодом правят в том же коммите, и на ревью видно, какое требование изменилось вместе с кодом.

    Таблица справа — этот файл. Строки шага лежат под текстом, их нужно перенести в таблицу. Документ один на весь разбор: на шагах про SLO и расчёты в нём появятся числа и новые свойства.

    Функциональное требование описывает действие системы. У каждого действия есть тот, кто его вызывает, поэтому в таблице для него отдельная колонка. Если вызывающего нет, это не требование, а преждевременное решение об устройстве.

    Строки пронумерованы, чтобы на них можно было ссылаться: «FR-5 требует, чтобы воркер умел прервать работу» короче пересказа.

    Очереди, пула и повторов в таблице нет. Это способы реализации, хотя их часто пытаются записать в функции первыми.

    В таблицу требований
    scraper/docs/REQUIREMENTS.md
    Функциональные0
    #Система умеетКто вызывает
    FR-1Перетащи сюда из текста слева
    FR-2Перетащи сюда из текста слева
    FR-3Перетащи сюда из текста слева
    FR-4Перетащи сюда из текста слева
    FR-5Перетащи сюда из текста слева
    FR-6Перетащи сюда из текста слева
    FR-7Перетащи сюда из текста слева
    FR-8Перетащи сюда из текста слева
    Нефункциональные0
    #СвойствоЦель
    Пока пусто
  5. Шаг 5 из 36

    Что вне рамок

    Список того, чего система не делает, стоит записать так же явно, как список функций. Через полгода кто-нибудь спросит, почему нет расписаний, и ответ должен быть в документе.

    • Обход капчи и антибот-защиты. Меняет и техническую, и юридическую сторону задачи. Если это нужно, это отдельный продукт со своими рисками.
    • Биллинг за выполненные джобы. Отдельный контекст со своим словарём. Если включить его сюда, придётся проектировать два сервиса одновременно.
    • Расписания. Сейчас джобу ставят снаружи. Для запуска по времени нужен отдельный планировщик: со своим состоянием, пропущенными запусками и правилом на случай, если предыдущий прогон ещё идёт.
    • Разбор ответа в структуру. Мы отдаём то, что вернула цель. Превращение HTML в таблицу — отдельная задача со своим циклом разработки, и сбор не должен от неё зависеть.

    Справа десять утверждений из того же разговора. Разложи их по видам. Часть свойств на слух похожа на функции, и на ревью обычно спорят как раз о них.

    Десять утверждений из разговора с продуктом. Разложи их: что система умеет, с какими свойствами она это делает и чего не делает вовсе.

    1. Принять джобу и сразу вернуть её идентификатор
    2. Показать по идентификатору, что с джобой происходит
    3. Отдать результат готовой джобы отдельным запросом
    4. Отменить джобу, которая ещё не выполнилась
    5. Приём отвечает быстрее 200 мс в 99 случаях из 100
    6. Принятая джоба не теряется при перезапуске воркера
    7. Один пользователь не занимает собой весь пул
    8. Приём работает, даже когда целевой API лежит
    9. Обходить капчу и антибот-защиту цели
    10. Считать деньги за выполненные джобы
    Разложено 0 из 10
  6. Шаг 6 из 36

    Нефункциональные требования

    Второй раздел — свойства. Нефункциональное требование не добавляет системе действий. Оно задаёт гарантию для тех, что уже есть.

    Такие требования удобно проверять через последствия, поэтому под каждым свойством написано, что сломается без него. Если написать нечего, это пожелание, и его лучше сразу вычеркнуть.

    Колонка «Цель» пока пустая. «Приём отвечает быстро» — свойство, «за 200 мс в 99 случаях из 100» — обещание. Обещания появятся на следующем шаге и встанут в эту колонку: у NFR-2 появится число, формулировка останется прежней.

    В таблицу требований
    scraper/docs/REQUIREMENTS.md
    Функциональные0
    #Система умеетКто вызывает
    Пока пусто
    Нефункциональные0
    #СвойствоЦель
    NFR-1Перетащи сюда из текста слева
    NFR-2Перетащи сюда из текста слева
    NFR-3Перетащи сюда из текста слева
    NFR-4Перетащи сюда из текста слева
    NFR-5Перетащи сюда из текста слева
    NFR-6Перетащи сюда из текста слева
    NFR-7Перетащи сюда из текста слева
  7. Шаг 7 из 36

    Числа, к которым привязываемся

    Свойство без числа нельзя проверить: о том, быстро ли отвечает приём, можно спорить сколько угодно. SLO задаёт долю запросов, порог и окно, за которое их считают.

    • Доступность приёма — 99,9% за 30 дней. Это 43 минуты в месяц, когда джобы не принимаются. Сто процентов обошлись бы намного дороже, чем стоит эта задача.
    • Приём отвечает за 200 мс в 99 случаях из 100. Считаем верхнюю границу, а не среднее: среднее скрывает медленные ответы, после которых клиент повторяет запрос.
    • 95% джоб стартуют в первую минуту. Это обещание про очередь. Оно нарушается, когда поток больше, чем успевает разобрать пул.
    • 99% принятых джоб доходят до конечного состояния без участия человека. Конечные состояния — «готово» и «упала окончательно». Джоба, которая навсегда осталась в «выполняется», хуже упавшей.

    По этим числам принимают решения. Нарушено первое — расширяем приём. Нарушено третье — увеличиваем пул или делим очередь.

    В таблицу требований
    scraper/docs/REQUIREMENTS.md
    Функциональные0
    #Система умеетКто вызывает
    Пока пусто
    Нефункциональные0
    #СвойствоЦель
    NFR-1Перетащи сюда из текста слева
    NFR-2Перетащи сюда из текста слева
    NFR-3Перетащи сюда из текста слева
    NFR-8Перетащи сюда из текста слева
  8. Шаг 8 из 36

    Расчёт: сколько джоб в секунду

    Пользователи и их привычки задают поток джоб. Поток и длительность джобы задают одновременность, а от неё зависит размер пула. Калькулятор справа: подставь свои значения.

    50 000 пользователей по 8 джоб в сутки дают 400 000 джоб в сутки, в среднем меньше пяти в секунду. Нагрузка неравномерная: утром и вечером людей больше, и при шестикратном пике получается 28 джоб в секунду.

    Число одновременно выполняемых джоб считается по закону Литтла: поток, умноженный на время одной джобы. 28 джоб в секунду по 20 секунд — 556 джоб одновременно. Размер пула определяется этим числом, а не количеством пользователей.

    Записями поток не исчерпывается. Пока джоба выполняется, о ней спрашивают, и спрашивают с той минуты, как она принята. Пять секунд в очереди плюс двадцать на выполнение при опросе раз в пять секунд — это пять запросов статуса и ещё один за результатом. Шесть чтений на одну запись, то есть 167 чтений в секунду на пике против 28 записей. Пул считается по записям, а база живёт под чтениями, и второе число в разборах теряют чаще первого. Что с ним делать — на шаге про кэш.

    Ожидание здесь ползунок, а не результат: каким оно будет, закон Литтла не говорит. Зато он переводит его в джобы — 139 штук в очереди на пике — и позволяет проверить NFR-8. Пока среднее ожидание пять секунд, обещание «95% джоб стартуют в первую минуту» выполнимо; подвинь ползунок к минуте, и оно нарушено, а заодно вырастет поток чтения.

    Последняя строка в таблице отвечает на вопрос, зачем очередь нужна вообще. Она считает другой случай: пул рассчитан по среднему потоку, и тогда за пять минут пика скапливается почти семь тысяч джоб. Разница между средним и пиком никуда не девается, она копится. Пул из расчёта по пику её не накапливает, но стоит дороже и простаивает большую часть суток.

    Попробуй изменить длительность. При двух секундах хватит пула на полсотни джоб и одной машины. При тридцати одновременно в работе больше восьмисот джоб: нужны кластер, автоскейл и отдельный бюджет. Задача та же, изменилась только скорость целевого API. Поэтому первым вопросом к продукту было, что именно мы скрейпим.

    В таблицу требований

    Что известно

    Что из этого следует

    Джоб в сутки
    400 000 шт.
    Средний поток
    4,63 джоб/с
    Поток на пике
    27,8 джоб/с
    Средний поток чтения
    27,8 чтений/с
    Чтение на пике
    167 чтений/с
    Одновременно в работе
    556 джоб
    Воркеров нужно
    28 шт.
    Ждёт очереди на пике
    139 джоб
    От приёма до результата
    25 с
    Очередь за 5 минут пика, если пул по среднему
    6 944 джоб
    Данных в сутки
    76,3 GB
    В хранилище
    2,24 TB

    Закон Литтла: 27,8 джоб/с × 20 с = 556 джоб в работе одновременно

    Опрос раз в 5 с — это 6 чтений на джобу: 167 чтений в секунду на пике против 27,8 записей.

    Нужен пул воркеров, который масштабируется отдельно от приёма. Без очереди между ними не обойтись.

  9. Шаг 9 из 36

    Расчёт: сколько воркеров

    Одновременность посчитана, но один воркер обрабатывает больше одной джобы. Скрейпер почти всё время ждёт сеть, и один процесс держит десятки открытых запросов. Если считать по джобе на воркер, машин понадобится в двадцать раз больше.

    При 20 джобах на воркер 556 одновременных джоб дают 28 воркеров. Сколько джоб реально держит один воркер, нужно измерить. Ограничением может стать память (ответ цели хранится в ней целиком), файловые дескрипторы, исходящие соединения или лимиты клиентской библиотеки.

    Из этого следуют два требования к архитектуре.

    Пул масштабируется по глубине очереди и возрасту самой старой джобы. Загрузка процессора у скрейпера низкая и в норме, и при аварии, поэтому автоскейл по ней ничего не заметит.

    Воркер может упасть в любой момент. Машина с 20 джобами в работе при падении оставляет 20 джоб, которые нужно вернуть в очередь. Как это сделать, разберём на шаге про жизненный цикл.

    В таблицу требований
    scraper/docs/REQUIREMENTS.md
    Функциональные0
    #Система умеетКто вызывает
    Пока пусто
    Нефункциональные0
    #СвойствоЦель
    NFR-9Перетащи сюда из текста слева
  10. Шаг 10 из 36

    Расчёт: сколько данных

    Результат не помещается в строку базы. 200 КБ на джобу при 400 000 джоб в сутки — это 76 ГБ в день и больше 2 ТБ за месяц хранения. Держать такой объём в таблице джоб нельзя.

    Из этой оценки следуют три решения.

    Результаты хранятся в объектном хранилище. Оно дешёвое, почти не ограничено по объёму и не зависит от размера ответа. В базе остаются метаданные: идентификатор, состояние, ссылка, размер, время. Запрос списка джоб не читает сами данные.

    Срок хранения — требование. «Храним месяц» значит, что кто-то удаляет старые данные. Иначе объём растёт линейно, и первым ограничением станет счёт за хранилище.

    Размер результата — входной параметр. В калькуляторе на шаге выше измени его с 200 КБ на 20 МБ: объём за месяц вырастет с 2 ТБ до 200 ТБ, и выбор хранилища станет вопросом бюджета.

    В таблицу требований
    scraper/docs/REQUIREMENTS.md
    Функциональные0
    #Система умеетКто вызывает
    Пока пусто
    Нефункциональные0
    #СвойствоЦель
    NFR-10Перетащи сюда из текста слева
  11. Шаг 11 из 36

    C1: система и всё, что с ней разговаривает

    Теперь можно рисовать. Первый кадр — C1: система одним блоком.

    Слева пользователи: человек в интерфейсе и чужой скрипт. Справа целевые API. Мы их не контролируем: они бывают медленными и недоступными и ограничивают нас лимитами. Вниз идёт вебхук для тех, кто попросил уведомление.

    Кадр показывает границу системы. Целевые API в неё не входят, хотя задача про них: их доступность не входит в наш SLO, лимиты устанавливают они, формат ответа тоже их. Мы можем только работать с ними как с ненадёжной внешней средой.

    Отсюда разделение обещаний. «Джоба принята» мы гарантируем сами. «Данные будут» — только если цель ответит.

  12. Шаг 12 из 36

    Вход: gateway

    Первый элемент внутри границы — точка входа, api-gateway. Сюда приходят все внешние запросы: и «поставь джобу», и «что с ней».

    Gateway проверяет, кто пришёл, отсекает лишний трафик и передаёт запрос дальше. Он выделен отдельно, потому что аутентификация и лимиты на клиента нужны и пишущей, и читающей стороне. Если реализовать их дважды, правила со временем разойдутся.

    • Новое на схемеapi-gatewayТочка входа: проверяет, кто пришёл, держит общие лимиты и раздаёт запрос внутрь. Про джобы не знает ничего.
  13. Шаг 13 из 36

    Запись: командная сторона

    За gateway стоит job-command — единственный сервис, который создаёт джобы, — и база jobs-pg, где они хранятся.

    Запись и чтение здесь разделены по CQRS: командная сторона создаёт и меняет джобы, читающая только отвечает на вопросы о них. Нагрузка у сторон разная. Джобу создают один раз, а спрашивают о ней десятки раз: человек обновляет список, скрипт опрашивает свою тысячу джоб. Если масштабировать обе стороны как одно целое, мощности придётся подбирать под чтение, и запись получит их с большим запасом. Подробнее о разделении — в главе про CQRS.

    Командный сервис проверяет параметры, проверяет ключ идемпотентности и записывает джобу.

    Очереди на кадре пока нет: сначала разберём запись, потом — как о ней узнают остальные части системы.

    • Новое на схемеjob-commandЕдинственный, кто создаёт джобу. В целевые API не ходит никогда.
    • Новое на схемеjobs-pgpostgresХранит джобы: чья, в каком состоянии, сколько было попыток.
  14. Шаг 14 из 36

    Как о джобе узнаёт очередь

    Простой вариант — записать джобу в базу и сразу отправить сообщение в очередь — ненадёжен. Если запись прошла, а публикация упала, джоба есть, но её никто не выполнит. Если сообщение ушло, а транзакция откатилась, очередь ссылается на несуществующую джобу. Одна транзакция не может охватить и базу, и брокер.

    Поэтому командный сервис пишет только в базу, а в шину событие отправляет отдельный публикатор job-publisher. Отдельная таблица для этого не нужна: джоба записывается с пустой колонкой published_at. На схеме появляются два элемента: публикатор и шина. Как это работает по шагам — дальше.

    • Новое на схемеjob-publisherНаходит в jobs строки с пустой отметкой публикации, отправляет событие в шину и проставляет время. Может отправить повторно — доставка «хотя бы раз».
    • Новое на схемеMQЗадания, которые ещё не взяли, и события о завершённых. Стоит внутри границы: шину мы эксплуатируем сами.
  15. Шаг 15 из 36

    Приём и публикация изнутри

    Приём — одна транзакция с одной записью: INSERT INTO jobs с пустым published_at. После COMMIT строка либо есть вместе с отметкой «опубликовать», либо её нет. Вторая таблица не нужна: состояние джобы и отметка о публикации лежат в одной строке.

    Публикатор работает отдельным циклом: выбирает строки с пустым published_at, отправляет событие в шину и проставляет время. Шина в транзакцию не входит, поэтому между отправкой и отметкой есть окно. Если публикатор упадёт в этом окне, после перезапуска он найдёт ту же строку и отправит событие ещё раз. Доставка «хотя бы раз» означает, что воркер может получить одну джобу дважды, поэтому взятие джобы в работу проходит через статус в базе.

    Тем же способом уходит и событие о завершении: воркер ставит статус «готово» вместе с пустым finish_published_at, а событие отправляет публикатор. Сам воркер в шину не пишет.

    Справа три сценария: всё прошло, сервис упал до COMMIT, публикатор упал после отправки. Сравни, что остаётся в таблице и в очереди.

    Сценарий
    Строк в jobs
    1
    published_at
    проставлен
    Сообщений в MQ
    1

    Джоба записана с пустым published_at. Публикатор нашёл её, отправил событие и проставил время. В очереди одно сообщение.

  16. Шаг 16 из 36

    Пул воркеров

    Пул воркеров scrape-worker выполняет работу: читает очередь, берёт джобу и обращается к чужому API.

    Пул масштабируется отдельно от приёма, по глубине очереди и возрасту самой старой джобы. По расчётам это 28 процессов по 20 джоб, и число зависит от скорости целевого API.

    Граница на кадре отделяет то, что мы контролируем, от того, что нет. Слева наш пул, справа чужие API. Со своей стороны мы можем задать таймаут, повторы и лимит на домен.

    • Новое на схемеscrape-workerБерёт джобу из очереди и обращается к чужому API. Больше никто в системе к внешним API не ходит.
  17. Шаг 17 из 36

    Взятие джобы и зависшие

    Получить сообщение из очереди ещё недостаточно, чтобы взять джобу. Воркер ставит ей статус «в обработке» и время взятия и только потом обращается к цели. Если придёт повторная копия сообщения, воркер увидит, что джоба уже занята, и пропустит её.

    У этого подхода есть известная проблема: воркер может упасть, не закончив работу. Тогда джоба останется «в обработке» навсегда, и дашборд её не покажет, потому что формально она выполняется.

    Решение использует то же время взятия. В сервисе планировщика job-scheduler есть крон-модуль: раз в минуту он находит джобы, которые слишком долго висят в обработке, ставит им статус queued и очищает published_at. Публикатор отправит их в очередь заново. На схеме C2 виден сам job-scheduler: крон — модуль внутри него, а модули покажем на уровне C3, после итоговой схемы C2.

    • Новое на схемеjob-schedulerПланировщик: выбирает следующую джобу, решает, когда попытки кончились, а его крон-модуль раз в минуту возвращает в очередь зависшие джобы.
  18. Шаг 18 из 36

    Чтение: статусы и история

    job-query закрывает требование из задания: видеть, что происходит с джобами.

    Читающая сторона отвечает на три вопроса: в каком состоянии джоба, какие джобы есть у пользователя и что происходило с конкретной джобой. Третий полезнее всего: история попыток с ответами цели показывает, например, «три таймаута, потом 401», и по ней понятно, что чинить.

    Gateway отправляет такие запросы сюда, а не в командный сервис. Требования к ним разные: приём должен отвечать за 200 мс и работать при падении целей, а список джоб может отдаваться на сотню миллисекунд медленнее.

    • Новое на схемеjob-queryОтвечает, в каком состоянии джоба и что с ней было. Масштабируется отдельно от записи: чтений больше.
  19. Шаг 19 из 36

    Кэш на читающей стороне

    167 чтений в секунду на пике против 28 записей, и все 167 идут в jobs-pg — в ту же базу, в которую пишет приём. Доступность приёма — это NFR-1, а значит, цикл опроса в чужом скрипте способен уронить приём джоб, к которому отношения не имеет.

    Между job-query и базой встаёт jobs-cacheпромежуточный буфер, в котором лежит ответ на GET /jobs/{id}: состояние, ссылка на результат, отметки времени. Ключ — job:{id}.

    Схема обращения — cache-aside. Читающая сторона сама идёт сначала в кэш: при попадании ответ отдаётся без обращения к базе, при промахе job-query читает строку из jobs-pg и сам кладёт её в кэш. Пишущая сторона о кэше не знает вовсе — ни job-command, ни воркер в него не пишут. Сквозная запись, при которой значение попадает в кэш в момент записи в базу, потребовала бы обратного: о кэше пришлось бы знать обоим.

    Списка джоб в кэше нет. У GET /jobs?status= ответ свой на каждое сочетание фильтра, сортировки и страницы, а недействительным его делает любая чужая джоба того же пользователя. Ключей много, коэффициент попаданий низкий, вытеснять приходится часто. В кэш идёт то, что запрашивают чаще всего и на что отвечают одинаково: одна джоба по идентификатору.

    Дальше — согласованность. Джоба перешла в конечное состояние, а в кэше лежит «выполняется»; клиент читает устаревшее значение и продолжает опрашивать. Нагрузку кэш в этом случае не снял, а задержку добавил. Поэтому у записи два срока жизни.

    Пока джоба не в конечном состоянии, TTL — две секунды. Одного опрашивающего такой TTL не спасает: при интервале опроса в пять секунд он промахивается каждый раз. Расчёт на то, что читателей у одной джобы обычно несколько — открытая вкладка, скрипт и обновление страницы дают несколько запросов в одну секунду, и до базы доходит один из них.

    После перехода в конечное состояние — до конца ретеншна. Состояние больше не изменится, а обращаются к такой джобе долго: за ссылкой на результат приходят и через час. Эти чтения до базы не доходят вовсе.

    Переключает между режимами событие. Факт о завершении джобы уже идёт в шине — его ради вебхука публикует job-publisher. Читающая сторона подписывается на ту же очередь и удаляет ключ; следующий запрос один раз сходит в базу и положит в кэш конечное состояние.

    Короткий TTL при этом остаётся, и не как дубль инвалидации. У cache-aside есть гонка, которую инвалидация по событию не закрывает: job-query промахнулся и читает строку из базы, в это время джоба завершается, событие приходит и удаляет ключ, которого ещё нет, — а следом job-query кладёт в кэш уже устаревшее «выполняется». Окно узкое, но ограничивает его только TTL.

    Вторая известная неприятность — лавина промахов, когда у горячего ключа истекает срок и в базу одновременно уходят все пришедшие за ним запросы. Здесь ей негде развернуться: при TTL в две секунды и десятке читателей на джобу это единицы запросов. На длинном TTL и более горячих ключах пришлось бы пускать к базе одного, а остальных заставлять ждать его ответа.

    Каким получится коэффициент попаданий, зависит от того, сколько читателей приходится на одну джобу. Это число измеряют — как и то, сколько джоб держит один воркер.

    В таблицу требований
    • Новое на схемеjobs-cacheredisПромежуточный буфер перед jobs-pg: ответ на GET /jobs/{id}. При попадании опрос статуса до базы не доходит.
  20. Шаг 20 из 36

    Как забрать результат

    Для клиента у джобы два состояния: она ещё выполняется или уже готова. Узнать это — обычный запрос в базу через job-query. В ответе статус и, если джоба готова, ссылка на результат.

    Сами данные лежат в results-s3: их объём — терабайты, и в базе состояний им не место. job-query не пропускает их через себя, а возвращает временную подписанную ссылку, и клиент скачивает результат из хранилища напрямую.

    Статус «готово» воркер ставит только после того, как файл уже в хранилище, поэтому по ссылке из готовой джобы результат всегда есть. Как устроено хранение и удаление по сроку — на шаге «Где лежат результаты».

    • Новое на схемеresults-s3s3Результаты джоб. Их объём — терабайты, поэтому они хранятся отдельно от базы состояний.
  21. Шаг 21 из 36

    C2 целиком

    Сервис собран целиком. Здесь только связи между сервисами; стрелки с двух концов — связь в обе стороны. Что и в каком порядке по ним идёт — на следующих двух шагах.

    Слева вход и две стороны, пишущая и читающая, которые масштабируются независимо; у читающей свой кэш, чтобы опрос статуса не доходил до базы. В центре база: состояние джоб с отметками публикации, благодаря которым запись и публикация не расходятся. В шину пишет только публикатор. Справа пул воркеров, единственный элемент, который обращается к чужим API, и хранилище результатов. Шина несёт две очереди: задания на выполнение и события о завершении — вторую слушает не только вебхук, но и кэш.

    На схеме пока нет трёх вещей, о которых пойдёт речь дальше: как воркер отличает временный отказ от постоянного, откуда берётся лимит на чужой домен и как пул делится между арендаторами.

  22. Шаг 22 из 36

    Запись по шагам

    Схема C2 показывает, кто с кем связан, но не в каком порядке. Здесь один запрос на постановку джобы, сверху вниз.

    Клиент получает 202 сразу после записи в jobs-pg — в этот момент джоба ещё никуда не отправлена. Дальше она идёт без него: публикатор находит строку и кладёт событие в шину, воркер берёт джобу, ходит в цель, складывает результат и ставит статус «готово».

  23. Шаг 23 из 36

    Чтение по шагам

    Запрос результата короче. Gateway спрашивает job-query, тот идёт сначала в jobs-cache и только на промахе — в jobs-pg. Если джоба готова, клиент получает подписанную ссылку и забирает данные из results-s3 сам, минуя наши сервисы.

    На кадре нарисован промах: он показывает всех участников. При попадании последовательность заканчивается на второй стрелке, и базы в ней нет вообще — а это, если считать по числам с шага про кэш, большинство запросов.

  24. Шаг 24 из 36

    C3: внутри job-scheduler

    Спускаемся на уровень ниже. C3 показывает, из чего состоит один сервис. Берём job-scheduler: на C2 он был одной коробкой, которая пишет в базу, а внутри у него три модуля.

    • stuck-sweeper — крон-модуль из шага про аренду. Раз в минуту находит джобы с истёкшей арендой, возвращает их в queued и очищает published_at.
    • retry-planner — решает, что делать после неудачной попытки. При временном отказе назначает следующую попытку с backoff, при постоянном отказе или исчерпанных попытках переводит джобу в failed. Подробнее — на шаге про повторы.
    • jobs-repository — единственный модуль, который обращается к jobs-pg. Два других работают через него.

    Ни один модуль не пишет в шину. Они меняют строку джобы и очищают отметку публикации, а событие отправляет job-publisher. Поэтому у планировщика та же гарантия, что у приёма: изменение состояния и отметка «опубликовать» — одна запись.

    • Новое на схемеstuck-sweeperКрон-модуль: раз в минуту возвращает в очередь джобы с истёкшей арендой.
    • Новое на схемеretry-plannerНазначает следующую попытку с backoff или переводит джобу в failed, когда попытки кончились.
    • Новое на схемеjobs-repositoryЕдинственный модуль, который обращается к jobs-pg.
  25. Шаг 25 из 36

    Контракт: расписка вместо результата

    POST /jobs отвечает 202 Accepted, а не 200 с данными: данных ещё нет, а открытое соединение оборвёт таймаут балансировщика. В ответе расписка — идентификатор и адрес, где следить.

    Ключ идемпотентности: клиент, не дождавшийся ответа, повторит запрос и получит ту же джобу, а не вторую.

    Результат отдаёт не API, а хранилище по подписанной ссылке: большие объёмы оно отдаёт дешевле. Вебхук дополняет опрос, но не заменяет: он может не дойти.

    • POST/jobs202 Accepted

      Поставить джобу. Ответ сразу, данных в нём нет: только расписка и адрес, где следить.

      запрос
      Idempotency-Keyurlparamspriority
      ответ
      idstatus: queuedlinks.self
    • GET/jobs/{id}200 OK

      Состояние джобы, число попыток и последняя ошибка цели.

      ответ
      statusattemptslast_error
    • GET/jobs?status=200 OK

      Свои джобы с фильтром по состоянию: скрипту с тысячей джоб хватает одного запроса.

      ответ
      items[]next_cursor
    • GET/jobs/{id}/result200 OK

      Не данные, а подписанная ссылка на хранилище со сроком действия.

      ответ
      urlexpires_at
    • POST{callback_url}200 OK

      исходящийМы зовём клиента, когда джоба закончилась. Может не дойти, поэтому опрос остаётся.

      запрос
      idstatusfinished_at
    В таблицу требований
    scraper/docs/REQUIREMENTS.md
    Функциональные0
    #Система умеетКто вызывает
    Пока пусто
    Нефункциональные0
    #СвойствоЦель
    NFR-5Перетащи сюда из текста слева
  26. Шаг 26 из 36

    Жизненный цикл джобы

    Состояний у джобы немного, и их видит пользователь.

    queuedrunningsucceeded или failed. Кроме того, cancelled: пока джоба не начала выполняться или пока воркер может её прервать. Из running есть переход обратно в queued: попытка не удалась, и джоба ждёт следующей.

    В running джобу переводит воркер, когда берёт её. В succeeded — тоже он, но только после того, как результат сохранён в хранилище, иначе появится «готово» без данных. В failed джоба переходит, когда закончились попытки. Это решает планировщик, а не воркер.

    Отдельно нужно обработать джобу, которую взяли и потеряли: воркер упал, а джоба осталась в running. Для этого используется аренда: воркер продлевает срок владения джобой, а если срок истёк, джоба возвращается в очередь. Без аренды джобы упавшей машины не вернутся, и дашборд этого не покажет, потому что формально они выполняются.

    • job-command
    • воркер
    • планировщик
    • клиент

    running

    Воркер взял джобу и держит аренду. Пока аренда продлевается, джоба его; истекла — возвращается в очередь.

    Куда дальше

    • succeededрезультат сохранён · воркер
    • queuedповтор после backoff, аренда истекла · планировщик
    • failedпопытки кончились · планировщик
    • cancelledотмена, воркер прерван · клиент
  27. Шаг 27 из 36

    Таймауты, повторы и что считать провалом

    Задание требует, чтобы джобы не падали без причины. Для этого нужно различать два вида отказов. Граница проходит не по коду ответа, а по тому, может ли повторная попытка дать другой результат.

    Временные отказы. Таймаут, разрыв соединения, 429, 503. Повторять стоит, но не сразу и не всем пулом одновременно. Нужен экспоненциальный backoff со случайным разбросом. Без разброса тысяча джоб, упавших в одну секунду, повторится тоже в одну секунду и снова перегрузит цель, которая только начала восстанавливаться.

    Постоянные отказы. 404, 401, неверные параметры, невалидный ответ. Повтор даст тот же результат и потратит ещё один запрос к цели. Такая джоба сразу переходит в failed.

    Число попыток ограничено. Джобы, которые исчерпали попытки, попадают в отстойник, где их разбирает человек. Например, сто джоб к одному домену с ответом 401 обычно означают один просроченный ключ.

    Повтор безопасен, только если цель не считает его новым действием. Для чтения это почти всегда так, но это стоит проверить до того, как включать повторы.

    В таблицу требований
    scraper/docs/REQUIREMENTS.md
    Функциональные0
    #Система умеетКто вызывает
    Пока пусто
    Нефункциональные0
    #СвойствоЦель
    NFR-11Перетащи сюда из текста слева
  28. Шаг 28 из 36

    Вежливость к чужим API

    Свой пул мы посчитали, но у цели свои ограничения. 28 воркеров по 20 джоб дают до 560 параллельных запросов. Если все они уйдут на один домен, нас заблокируют раньше, чем закончатся наши мощности.

    Поэтому лимит задаётся для каждой цели: столько-то запросов в секунду и столько-то одновременных соединений на домен. Это токен-бакет на хост, общий для всех воркеров, поэтому он хранится не в памяти процесса, а там, где его видят все воркеры.

    Из этого следует, что очередь делится по целям, а не по пользователям. В общей очереди джоба к свободному домену ждёт за джобами к занятому, и пул простаивает. С очередью на домен занятый домен задерживает только свои джобы.

    Заголовок Retry-After в ответе цели нужно соблюдать. Если его игнорировать, временное ограничение может превратиться в постоянную блокировку.

    В таблицу требований
    scraper/docs/REQUIREMENTS.md
    Функциональные0
    #Система умеетКто вызывает
    Пока пусто
    Нефункциональные0
    #СвойствоЦель
    NFR-6Перетащи сюда из текста слева
  29. Шаг 29 из 36

    Шумный сосед

    Изоляция арендаторов записана в требованиях одной строкой. Симулятор справа показывает, что она означает на практике.

    Три пользователя: магазин и аналитик с небольшим потоком и краулер, который сразу поставил 400 джоб. Пул из восьми воркеров по четыре секунды на джобу обрабатывает две джобы в секунду. Включи общую очередь: за первые секунды краулер выполняет десятки джоб, а у соседей нет ни одной готовой, и их очередь растёт. Система при этом работает правильно: она выполняет джобы в порядке поступления.

    Переключи порядок на «поровну». Воркеры, поток и длительность те же, изменилось только правило выбора следующей джобы. Соседи получают свои результаты, краулер получает треть пула, и его очередь разбирается дольше.

    Справедливый порядок усложняет планировщик: нужна очередь на каждого арендатора, учёт выданных слотов и правило для арендатора, который долго молчал и вдруг пришёл с работой. Следующий шаг — веса: платящим клиентам можно выделять больше, чем бесплатным. Размер весов решает продукт.

    Порядок00:00

    Запусти и смотри на цвета слотов в пуле — по ним видно, чья работа его занимает.

  30. Шаг 30 из 36

    Где лежат результаты

    Воркер сохраняет результат в объектное хранилище. Причин две: объём, посчитанный на шаге про данные (больше 2 ТБ в месяц), и назначение базы состояний — отвечать на запросы о состоянии, а не хранить большие файлы.

    Порядок записи: сначала данные в хранилище, потом ссылка в базу, потом состояние succeeded. При обратном порядке клиент увидит «готово» раньше, чем данные станут доступны, и не найдёт результат.

    Пользователь получает временную подписанную ссылку. Данные скачиваются из хранилища напрямую, минуя наш API, а ссылка перестаёт работать по истечении срока.

    Удаление по сроку хранения тоже нужно продумать. Хранилище умеет удалять файлы по правилу, но база об этом не узнает, и джоба останется succeeded со ссылкой на удалённый файл. Поэтому нужно либо показывать срок хранения в ответе, либо переводить джобу в отдельное состояние после удаления данных.

  31. Шаг 31 из 36

    Что мерить

    Метрик у такой системы много, но для решений достаточно четырёх.

    • Возраст самой старой джобы в очереди. Важнее длины очереди: тысяча джоб, которые разбираются за минуту, — нормальная работа, а десять джоб, которые ждут час, — авария. По этой же метрике удобно масштабировать пул.
    • Время до первой попытки, p95. Прямо соответствует SLO. Остальные метрики помогают понять, почему он нарушен.
    • Доля отказов по классам и доменам. Общая доля отказов смешивает просроченный ключ одного клиента с проблемами одной цели. Сто ответов 401 от одного домена — это один инцидент.
    • Глубина очереди по арендаторам. Показывает шумного соседа до того, как на него пожалуются.

    Отдельно нужно отслеживать джобы, взятые упавшим воркером. Формально они в running, и обычные метрики считают их нормальными. Их находят только по времени в состоянии: если джоба выполняется дольше срока аренды, она должна попадать в отчёт.

    В таблицу требований
    scraper/docs/REQUIREMENTS.md
    Функциональные0
    #Система умеетКто вызывает
    Пока пусто
    Нефункциональные0
    #СвойствоЦель
    NFR-7Перетащи сюда из текста слева
  32. Шаг 32 из 36

    Новое требование: приоритеты

    Через полгода продукт приходит с новой просьбой: часть джоб срочная. Магазину нужна цена конкурента до конца распродажи, а ночная выгрузка аналитика может подождать.

    В POST /jobs появляется поле priority: high, normal или low, по умолчанию normal. Оно пишется в строку jobs вместе с остальными полями джобы.

    Одна очередь приоритета не даёт: брокер отдаёт сообщения по порядку. Поэтому очередей работы теперь три: jobs.high, jobs.normal и jobs.low. Публикатор читает приоритет из строки и кладёт событие в нужную. Outbox не меняется, меняется только маршрут. Повтор после сбоя возвращается в ту же очередь.

    Главное здесь — как воркер выбирает очередь. Если всегда сначала high, то при постоянном потоке срочных джоб low не стартует никогда. Поэтому воркер берёт с весами: из десяти джоб шесть из high, три из normal и одну из low. Пустая очередь отдаёт свою долю остальным.

    Справедливость по арендаторам остаётся внутри каждой очереди: приоритет решает, из какой очереди брать, а справедливость — чью джобу из неё.

    В таблицу требований
    • Новое на схемеjobs.highСрочные джобы. Воркер берёт отсюда шесть из десяти.
    • Новое на схемеjobs.normalПриоритет по умолчанию. Три из десяти.
    • Новое на схемеjobs.lowТо, что может подождать. Одна из десяти, но никогда не ноль.
  33. Шаг 33 из 36

    Приоритеты в симуляторе

    Тот же пул из восьми воркеров и тот же краулер с залпом в 400 джоб. Теперь у каждого арендатора свой приоритет: магазин — high, краулер — normal, аналитик — low.

    Запусти на «строго по приоритету». Магазин получает слот сразу, краулер разбирает свой залп, а аналитик не стартует ни разу, пока у краулера есть работа. Никто не ошибся: low просто всегда последний.

    Переключи на веса. Краулер и магазин по-прежнему впереди, но аналитик получает свою долю — одну джобу из десяти, — и его очередь движется.

    Поменяй приоритеты над схемой: поставь краулеру high. Залп всё равно не займёт весь пул: у high шесть взятий из десяти, остальное достаётся другим очередям.

    Порядок00:00
    Магазин
    Аналитик
    Краулер с залпом

    Запусти и смотри на цвета слотов в пуле — по ним видно, чья работа его занимает.

  34. Шаг 34 из 36

    Какую MQ взять

    До этого шага MQ была коробкой без имени. Так и нужно, пока не ясно, что от неё требуется. Теперь требования записаны, и выбирать можно по ним, а не по привычке.

    Нам нужна очередь работы: подтверждение каждой джобы, отложенный повтор, три очереди по приоритету и очередь для того, что не разобралось. Плюс поток «джоба закончилась», который читают вебхук и кэш. Поток небольшой: 28 джоб в секунду на пике выдержит любой брокер.

    Kafka отпадает первой. Это журнал, а не очередь: она помнит смещение в партиции, а не судьбу каждого сообщения. Медленная джоба держит всю партицию, отложенный повтор и DLQ приходится собирать руками. Для потока событий Kafka хороша, для раздачи работы — нет.

    SQS закрывает очередь работы целиком, но привязывает к AWS, а для раздачи одного события нескольким потребителям нужен ещё SNS. Если сервис и так живёт в AWS, это разумный выбор.

    RabbitMQ и NATS JetStream набирают поровну. Берём RabbitMQ: приоритеты и dead-letter у него из коробки, а отложенный повтор собирается из TTL и dead-letter. NATS проще в эксплуатации и сам умеет задерживать повтор. Если команда уже его знает, выбор можно развернуть.

    Что нужно намKafkaRabbitMQберёмNATS JetStreamSQS
    NFR-3Подтверждение каждого сообщениятолько смещение в партиции: медленное сообщение держит остальныеack и nack на каждое сообщение, возврат при обрывеack на сообщение, повтор после AckWaitvisibility timeout: не удалил — сообщение вернётся
    NFR-11Отложенный повтор (backoff)нет: retry-топики собираются руками~TTL и dead-letter или плагин задержкиNakWithDelay и BackOff в настройках консьюмераChangeMessageVisibility, задержка до 15 минут
    NFR-12Очереди по приоритету~топик на приоритет, веса — в своём кодеочередь на приоритет или x-max-priority~subject на приоритет, веса — в своём коде~очередь на приоритет, веса — в своём коде
    NFR-3Очередь для неразобранного (DLQ)нет: DLQ-топик пишет сам потребительdead-letter exchange из коробки~MaxDeliver и событие о сдаче, очередь — свояredrive policy из коробки
    FR-7Одно событие — нескольким потребителямгруппы потребителей читают один топикfanout-exchangeнесколько консьюмеров на один потокнет: перед очередями нужен SNS
    NFR-9Пик 28 джоб/сс запасом на порядкидесятки тысяч в секундудесятки тысяч в секундупочти без предела
    Цена эксплуатациикластер, партиции, ребалансы~кластер и quorum-очередиодин бинарник, кластер собирается простоуправляет AWS
    Без привязки к облакуоткрытый, есть у всех облаковоткрытыйоткрытыйтолько AWS
    Итого3.5 / 87 / 87 / 85.5 / 8

    есть из коробки~ можно, но руками нет

  35. Шаг 35 из 36

    Чем заплатили

    Система выполняет все четыре обещания из задания. Вот что за это пришлось отдать.

    Асинхронность усложнила работу клиенту. Он не получает данные в ответ на запрос: хранит идентификатор, опрашивает сервис или ждёт вебхук и обрабатывает состояние failed. Для скрипта это привычно, для страницы в браузере это дополнительный код, который не понадобился бы с синхронным API.

    Очередь — ещё один компонент, который может отказать. Она сглаживает пики, но сама может упасть, переполниться или потерять сообщения при неверной настройке. Её нужно мониторить.

    Справедливость усложнила планировщик. Вместо «взять следующую» нужны очереди на арендаторов, учёт и веса. Этот код нужно тестировать и объяснять новым людям в команде.

    Результаты хранятся в двух местах. Метаданные в базе, данные в хранилище, и между ними возможно расхождение. Поэтому порядок записи пришлось оговорить отдельно.

    Решение придётся пересмотреть, если появится требование получать данные сразу (для быстрых целей понадобится отдельный синхронный путь) или если поток вырастет настолько, что одна очередь перестанет справляться и её придётся делить по доменам физически.

  36. Шаг 36 из 36

    Ответ: решения против требований

    В начале разбора мы записали требования, по ходу к ним добавились числа. Осталось проверить, что на каждое есть ответ.

    Слева решения, принятые по ходу разбора, справа таблица требований целиком. Положи каждую карточку на строку, которую она закрывает. У некоторых решений строк несколько, у некоторых строк — несколько решений.

    Если карточка легла не туда, под ней будет написано, на что она отвечает на самом деле. Требование, которое останется без карточки, означает, что разбор на него не ответил и схему нужно дорабатывать.

    Решения

    Закрыто 0 из 23

    Функциональные

    • FR-1Принять джобу — что забрать, откуда, с какими параметрами — и вернуть идентификатор
    • FR-2Показать состояние джобы по идентификатору
    • FR-3Показать список своих джоб с фильтром по состоянию
    • FR-4Отдать результат готовой джобы
    • FR-5Отменить джобу, которая ещё не выполнилась
    • FR-6Повторить джобу — новой джобой со ссылкой на прежнюю
    • FR-7Позвать обратно, когда джоба закончилась
    • FR-8Показать историю попыток джобы с ответами цели
    • FR-9Задать джобе приоритет: high, normal или low

    Нефункциональные

    • NFR-1Приём переживает падение целей
    • NFR-2Приём отвечает быстро
    • NFR-3Принятая джоба не теряется
    • NFR-4Арендаторы изолированы
    • NFR-5Повтор приёма не создаёт вторую джобу
    • NFR-6Мы вежливы к целям
    • NFR-7По джобе видно, что с ней было
    • NFR-8Джоба стартует без долгого ожидания
    • NFR-9Система выдерживает пиковый поток
    • NFR-10Результат хранится срок ретеншна
    • NFR-11Временные отказы цели не роняют джобу
    • NFR-13Чтение статуса не нагружает базу приёма
    • NFR-14Готовая джоба видна сразу
    • NFR-12Срочные джобы идут быстрее, низкий приоритет не голодает

Листать шаги можно стрелками ← и →.