← Все проекты уровня 15
Уровень 15 · Платформа выполнения задачВариант B

Планировщик обработки данных

DAG-планировщик с локальными источниками и приёмниками, сохранением прогресса и безопасным повтором этапов.

Техническое задание

Pipeline и данные

Создайте платформу запуска локальных ETL pipeline. План обработки состоит из 1–30 шагов DAG с именем, типом ingest, transform или export, конфигурацией, зависимостями и политика повторов: максимум 3 попытки шага, первая сразу и ещё две через 1 и 2 секунды; ошибка входных данных не повторяется. Коннекторы читают локальные CSV/JSON и PostgreSQL sandbox в Compose; облака, платные сервисы и произвольные команды не нужны. Ingest читает fixture из разрешённого каталога до 50 MiB; CSV-разбор обрабатывает quoted multiline поля как одну логическую запись, границы partition проходят только между CSV-записями. Transform выбирает из фиксированного набора операций: trim/normalize whitespace, map column, filter by literal equality, parse date ISO, aggregate count/sum. Числа — int64; overflow завершает шаг ошибкой; пустой набор даёт count=0, sum=0, min=null, max=null. Экспорт записывает результат в локальный каталог либо в sandbox таблицу с ключом run_id/step_id/record_key. Типы колонок и ошибки строк видны; режим pipeline задаёт остановку либо запись ошибок в reject-файл.

REST и расписание

Локальный seed создаёт pipeline_admin с токеном в Authorization header; только эта роль меняет pipeline и расписания. POST /api/v1/pipelines создаёт pipeline, GET/PATCH /pipelines/{id} читает и меняет неактивную версию, POST /pipelines/{id}/runs запускает вручную, GET /runs/{id} показывает шаги, POST /runs/{id}/cancel отменяет, POST /pipelines/{id}/schedules задаёт расписание UTC, DELETE отменяет расписание. Изменение плана создаёт revision; запущенный run использует её снимок. Расписание задаётся интервалом в минутах от 5 до 1440 или daily UTC, без cron-выражения. Pipeline запускает один run; следующие три срабатывания ставятся в очередь. Четвёртый триггер и все последующие получают skipped с причиной queue_full; уже поставленные в очередь запуски не удаляются и выполняются по порядку. Run имеет ID и trigger key; повтор ключа не создаёт новый run.

Контрольная позиция и повтор

Координатор хранит запуски, ревизии, состояния этапов, input_digest и output_digest, контрольные позиции, попытки, сроки прав обработки, счётчики отклонённых строк и аудит в PostgreSQL с миграциями и ручным SQL. Экспорт использует уникальный ключ run_id/step_id/record_key и upsert; при восстановлении запись не удваивается. Новый run имеет новый run_id и считается новым экспортом. Если шаг упал, только этот шаг и зависящие от него запускаются заново; успешно завершённая независимая ветка сохраняется. Политика on_failure плана обработки равна stop_dependents или continue_independent; всегда применяется одна политика, записанная в revision. Срок права обработки исполнителя — 30 секунд, heartbeat — каждые 10 секунд, одновременно работают до двух исполнителей. Timeout шага 5 минут, партия 500 записей, память исполнителя максимум 512 MiB. Очередь ограничена 20 runs.

Структура программы

Координатор и исполнители запускаются раздельно; ядро проверяет DAG и преобразует данные, транспорт, хранилище и запуск шагов не входят в него. main.go только загружает конфигурацию, связывает зависимости и управляет запуском. Пакеты без циклов, общих utils и интерфейсов без потребителя. README показывает зависимости; тесты ядра без HTTP/БД, интеграционные отдельно.

Приёмка

Compose включает координатор, два исполнителя и PostgreSQL. Локальные тестовые данные тестируют DAG, ошибки разбора, файл отклонённых записей, повтор расписания, остановку после экспорта и до фиксации контрольной позиции и отказ ветки. После перезапуска экспорт не содержит дубликатов; незавершённая порция выполняется повторно. Проверяйте ограниченную очередь при пиковой подаче 30 запусков: не более 20 queued/running суммарно; оставшиеся trigger получают skipped/busy результат. Критерии готовности: повтор запрос запуска создаёт одну запись, зафиксированная ревизия не меняется во время run, результат текущего шага невидим до фиксации, а отказ влияет только на downstream и выбранную политику. Отчёт воспроизводит импорт 100 000 CSV-записей, отказ одного исполнителя посередине этапа и продолжение после перезапуска; он содержит команды, идентификаторы трассировки, длительность и пиковую память среды.

Критерии готовности

Ожидаемый результат

Проверяемый результат

  • README показывает границы и направление зависимостей; main.go только связывает конфигурацию и запуск; модульные тесты ядра обходятся без HTTP и БД, интеграционные тесты отделены.
  • Docker Compose запускает координатор, два исполнителя и PostgreSQL; pipeline импортирует локальный CSV, преобразует данные и экспортирует результат.
  • API сохраняет revision pipeline, расписания UTC, trigger key, trace_id, статусы этапов и причину пропуска перегруженного расписания.
  • Контрольная позиция и результат фиксируются безопасно; остановка между экспортом и фиксацией контрольной позиции не удваивает строки.
  • Тесты проверяют выбранную политику отказа, файл отклонённых записей, повтор run, новую revision, расписание, burst и рестарт.
  • Один pipeline имеет не более одного активного run и трёх ожидающих; общая очередь ограничена 20 runs.
  • Каждый ответ содержит trace_id; журналы коррелируют его с pipeline_id, run_id, step_id, attempt_id и worker_id. Метрики показывают очередь, возраст права обработки, retries, rejected rows и пиковую память.
  • В отчёте воспроизводится обработка 100 000 CSV строк с временем и памятью на указанной конфигурации; вывод не выдаётся за универсальный SLA.