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-записей, отказ одного исполнителя посередине этапа и продолжение после перезапуска; он содержит команды, идентификаторы трассировки, длительность и пиковую память среды.