Как DVT выполняет проект

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

Задача попадает в очередь, затем её забирает свободный воркер

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

Задача, очередь и воркер

Запуск вручную, по расписанию или через API создаёт задачу и помещает её в общую очередь. Свободный воркер забирает задачу и выполняет проект. Один воркер выполняет одну задачу за раз и не набирает задачи про запас. Несколько воркеров позволяют выполнять несколько проектов одновременно.

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

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

Состояние

Что происходит

QUEUED

Задача ожидает свободного воркера

STARTED

Воркер принял задачу и подготавливает выполнение

RUNNING

Выполняются ноды проекта

CANCEL_REQUESTED

Запрошена остановка, выполнение ещё завершается

SUCCESS

Выполнение завершилось успешно

ERROR

Ошибка ноды или аварийное прерывание

CANCELLED

Задача отменена; причину нужно уточнить в её сообщении

Если у одного проекта накопилось несколько ещё не начатых обычных запусков, более новый может заменить старый в очереди. Пересчёт метаданных после правки графа также может вытеснить ожидающий полный запуск. Это правило не отменяет уже выполняющуюся задачу и не заменяет запуски подпроектов через Execute Project.

История проекта: выполняющаяся задача Running и завершённые задачи Success

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

Как выполняются ноды

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

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

Подробнее: подключения нод и обработка ошибок в пайплайне.

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

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

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

Ленивые вычисления

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

Нода результата — например, запись в БД или сохранение в Parquet, CSV, Excel — запускает накопленный план от начала до конца.

Чтение и преобразования строят план, нода результата выполняет его

Что видно в интерфейсе

Как это понимать

Чтение быстро завершилось, запись работает долго

Во время записи выполняются также чтение и преобразования

Индикаторы чтения и преобразований активны во время записи

Через эти шаги реально проходят данные

Ошибка преобразования появилась на ноде результата

Ошибка обнаружена при вычислении плана на реальных значениях

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

При обработке частями данные проходят через шаги без обязательного сохранения всего набора

Например, в цепочке «БД → фильтр → новая колонка → запись» первые шаги могут быстро подготовить план. Когда начинается запись, DVT читает данные, применяет фильтр и вычисляет новую колонку.

Что происходит до выполнения плана

Read Table DB V3 и Read Query DB V3 уже на своём шаге обращаются к базе: получают структуру, оценивают размер строки, количество строк и диапазон колонки разбиения, определяют границы частей. На большой таблице эти служебные запросы могут занять заметное время, особенно без индекса. SQL в Read Query может выполняться несколько раз при разметке и затем при чтении партиций.

Источник

Когда читает данные

Read Table DB V3, Read Query DB V3

Лениво, отдельными запросами по партициям

Load CSV, Load Parquet, Load Excel

При выполнении подготовленного плана

Read Queue Topic

Читает сообщения на собственном шаге, затем передаёт результат дальше

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

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

Служебные запросы при подготовке

Ленивые вычисления не исключают обращений к источнику до основной выгрузки. У чтения метаданных, построения Dask-датафрейма и расчёта партиций разные задачи и ограничения.

Операция

Read Table DB V3

Read Query DB V3

Отдельное получение метаданных

Структура и комментарии через SQLAlchemy Inspector

LIMIT 0, TOP 0 или WHERE 1=0; для ClickHouse — DESCRIBE

Подготовка метаданных Dask-датафрейма

Выборка максимум одной строки

Выборка максимум одной строки: LIMIT 1 или аналог СУБД

Оценка размера строки для автоматического разбиения

Возможна выборка до 1000 строк

Возможна выборка до 1000 строк

Основная выгрузка

Чтение партиций согласно плану

Выполнение запросов для партиций согласно плану

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

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

Партиции и потоки

Потоки обрабатывают разные части данных

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

Настройка

На что влияет

Число воркеров

Количество одновременно выполняющихся проектов

Количество потоков проекта

Количество одновременно обрабатываемых частей данных внутри проекта

Разбиение на партиции

Размер и количество частей, распределяемых между потоками

Связи «Сигнал»

Порядок выполнения связанных нод

В Read Table колонку разбиения можно указать явно; иначе используется первичный ключ, если он состоит из одной колонки. Без подходящей колонки нода завершается ошибкой. В Read Query колонку разбиения нужно задать обязательно.

Автоматический расчёт учитывает размер строки, целевой размер части и ограничения числа строк. В описанной конфигурации ориентир составляет около 29 МБ: 16 МБ × коэффициент 1,8. Минимум — обычно 10 000 строк, максимум зависит от СУБД; в руководстве для PostgreSQL указан предел 80 000 строк. Всего частей — не больше 5000. Эти параметры уточняйте у администратора.

Для упорядочиваемой колонки без пустых значений используются интервалы значений с примерно одинаковым числом строк. При пустых значениях или неподходящем типе применяется разбиение по хэшу. Малое число разных значений может привести к неравномерным частям: одинаковые значения оказываются вместе.

Прогресс ноды растёт по мере завершения частей. Асинхронное исполнение нод и обработка партиций в нескольких потоках — разные механизмы.

Для 450 партиций по 2 секунды каждая идеальная оценка при 8 потоках — 450 × 2 / 8 ≈ 113 секунд; при одном потоке — 15 минут. Реальное время зависит от базы, сети, преобразований и равномерности распределения строк между частями.

Почему больше потоков не всегда быстрее

Потоки особенно полезны, когда основное время уходит на ожидание сети, дисков и ответа БД. При перегруженном источнике или приёмнике увеличение числа потоков повышает очередь запросов. Сложный построчный Python-код может упираться в процессор и почти не ускоряться.

Каждый поток использует соединение с БД. При 16 потоках источник может получать до 16 одновременных запросов, а приёмник — до 16 потоков вставки. Учитывайте также нагрузку других проектов.

Увеличивайте число потоков постепенно и измеряйте время, память и нагрузку на базы. Диапазон 4–16 можно использовать для пробных измерений; подходящее значение зависит от проекта и ресурсов. Для SQL-баз начинайте с небольшого числа потоков и учитывайте предел соединений.

Нода результата

Управление параллельностью по руководству программиста

Write DataFrame To DB V3 / V4

«Количество потоков» проекта; в описанной конфигурации по умолчанию 8, верхний предел 32

Save Parquet

2 исполнителя; для FTP, SFTP, SMB — 1; настройка потоков проекта не влияет

Save CSV

Стандартная параллельность Dask по числу ядер воркера

Save Excel, преобразование в JSON

Вся таблица собирается в памяти перед сохранением

Для SQL-баз в описанной конфигурации пул одной ноды допускает до 15 соединений, ожидание свободного соединения — до 30 секунд. При превышении ожидания возможна ошибка QueuePool/timeout. Поэтому увеличение потоков выше 15 обычно не ускоряет такую запись. Для ClickHouse число соединений записи определяется числом потоков. Чтение и запись имеют независимые соединения; учитывайте совокупную нагрузку на обе базы.

Старые Write DataFrame To DB и V2 в руководстве обозначены как устаревшие; используйте V3 или V4.

Память и операции с полным набором

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

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

Для обработки партициями можно использовать приблизительную оценку:

Память ≈ количество потоков × размер партиции × 3 + 0,5 ГБ.

Множитель 3 учитывает промежуточные копии, а 0,5 ГБ — базовые нужды рабочего процесса и буфер кэша. Это оценка для планирования ресурсов, а не гарантия потребления или верхняя граница.

Потоков

Партиция 16 МБ

Партиция 100 МБ

4

≈ 0,7 ГБ

≈ 1,7 ГБ

8

≈ 0,9 ГБ

≈ 2,9 ГБ

16

≈ 1,3 ГБ

≈ 5,3 ГБ

32

≈ 2,0 ГБ

≈ 10 ГБ

Минимальные требования к оборудованию описаны отдельно.

Операции с повышенным расходом памяти

Save Excel и преобразование таблицы в JSON собирают весь набор в памяти. Сортировка, установка индекса, pivot, join и группировка могут перераспределять данные между партициями и заметно увеличивать потребление памяти. Фильтруйте строки и выбирайте нужные колонки до таких операций.

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

Предел 32 в описанной конфигурации относится к числу потоков, а не к объёму памяти. Конкретные параметры обработки и защиты уточняйте у администратора.

Несколько нод результата

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

Кэширование включается отдельно для ноды в контекстном меню; оно доступно у нод, выдающих таблицу или JSON. Время хранения задаётся настройками проекта «Включить TTL кэш» и «Время жизни (сек)». Согласно руководству, при значении 0 используется срок 10 минут.

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

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

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

Подробнее: кэш и повторные запуски.

Запись результата

В режимах truncate и upsert ноды Write DataFrame to DB V4 потоки сначала записывают данные во временную таблицу. Когда подготовлен весь слепок, он переносится в рабочую таблицу. В append пачки сразу добавляются в рабочую таблицу.

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

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

Write DataFrame To DB V4 позволяет создать целевую таблицу на этапе «Таблица» и настроить её схему через SQL или конструктор. Создание таблицы и режим заполнения данными — разные настройки. Колонки сопоставляются по именам; правила для лишних и недостающих колонок и ручное сопоставление задаются в ноде.

См. режимы записи и ошибки.

Остановка и повторный запуск

Обычная остановка

  1. Откройте выполняющийся проект.

  2. Нажмите команду остановки выполнения.

  3. Дождитесь окончательного статуса.

Кнопка «Остановить» на правой панели редактора

Кнопка «Остановить» — красная кнопка с квадратом. Hard Stop на этом кадре не показан.

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

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

Принудительная остановка

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

Остановка не откатывает автоматически уже записанные в БД или файлы данные. Перед повтором проверьте журнал и фактический результат.

Аварийно прерванная задача сама по себе не перезапускается: она могла частично записать данные. Повторы после ошибки для проектов по расписанию настраиваются отдельно.

Если запуск завершился ошибкой

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

Причина завершения

Что проверить

Ошибка SQL, Python или параметров ноды

Ноду и сообщение в журнале

Недоступен источник

Подключение и доступность источника

WORKER_LOST

Доступность воркера и уже записанные данные; передать администратору время запуска

OOM_GUARD

Объём обработки, размеры партиций, число потоков и тяжёлые операции

Недостаточно обработчиков для ожидания подпроекта

Доступность воркеров и настройки Execute Project

Execute Project с ожиданием результата удерживает воркер родительской задачи; подпроекту нужен ещё один доступный исполнитель. При одном воркере такой запуск отклоняется. Потеря необходимой ёмкости во время ожидания также может завершить родительскую задачу ошибкой.

С чего начать настройку

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

Ситуация

Следующее действие

Не хватает памяти

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

Выгрузка медленная, ресурсы свободны

Постепенно увеличьте число потоков; проверьте индекс колонки сегментации и размер пачки записи

База источника перегружена

Уменьшите число потоков и разнесите расписания проектов

Одна часть заметно медленнее остальных

Проверьте распределение значений колонки сегментации

Маленькая таблица разбита неудобно

Задайте количество партиций вручную

Настройки разбиения описаны в Read Table from DB V3 и Read Query from DB V3.