Единый декларативный контракт
В мире, где количество пайплайнов исчисляется сотнями, копипаст становится главным врагом инженера. Каждая новая таблица не должна требовать нового DAG-файла. Наш фреймворк заменяет императивный код на строго типизированную структуру BaseTransferConfig.
Разработчик больше не пишет логику чтения и записи. Он лишь декларирует намерения: Откуда (From) и Куда (To). Вся остальная магия — от обработки ошибок до оптимизации батчей — скрыта за абстракцией.
Типизация и Валидация
Использование Pydantic-датаклассов гарантирует, что конфигурация пайплайна валидна еще до его запуска. Ошибки «опечаток» в названиях полей исключены на этапе компиляции.
Универсальный Интерфейс
Поддержка Postgres, ClickHouse, Kafka, Databus и OpenSearch через единый интерфейс BaseTransportInterface. Смена приемника данных занимает одну строку в конфиге.
Zero-Config Metadata: Рефлексия источника
Фреймворк обладает «зрением». Ему не нужно вручную перечислять колонки или искать первичные ключи. Модули PostgresMetadata и KafkaMetadata выполняют глубокую рефлексию системных каталогов СУБД.
- Авто-детект ключей: Анализ
pg_attributeиsystem.columnsдля поиска Primary Key, Auto Increment или индексов сортировки. - Структурная целостность: Автоматическое получение списка колонок, их типов, Nullable-статуса и даже комментариев к полям.
SELECT a.attname, format_type(a.atttypid, a.atttypmod)
FROM pg_attribute a WHERE a.attrelid = 'table'::regclass
Динамический маппинг типов и авто-DDL
При первом обнаружении новой таблицы фреймворк выступает в роли архитектора. Он транслирует типы данных из источника (например, Postgres) в целевую систему (ClickHouse) и автоматически выполняет DDL-запрос.
Создаются отказоустойчивые таблицы ReplicatedReplacingMergeTree с правильным ORDER BY и сохранением бизнес-комментариев. Это позволяет разворачивать новые потоки данных за секунды.
Точность Маппинга
Сложные типы, такие как timestamp with time zone, автоматически конвертируются в DateTime64(6), а массивы — в Array(T).
Гибкость Схемы
Поддержка input_format_skip_unknown_fields позволяет эволюционировать схеме источника без падения пайплайна.
Параллельные потоки (Parallel Flow)
Чтение больших таблиц одним потоком — это узкое горлышко. Фреймворк реализует математическое секционирование данных на стороне источника, позволяя распараллелить нагрузку на N воркеров без конфликтов.
// Равномерное распределение нагрузки по хэш-ключу
Для ClickHouse используется более эффективная функция modulo(cityHash64(...)). Это обеспечивает линейное масштабирование скорости забора данных при увеличении количества воркеров.
Мультитранспортность: Single Source, Multi-Sink
Главная архитектурная особенность — принцип «одного запроса». Мы один раз вычитываем пакет данных из мастер-базы и параллельно доставляем его во все назначенные системы. Это кратно снижает IOPS источника.
ClickHouse Native
Вставка через JSONEachRow с контролем размера батча не в строках, а в байтах. Это предотвращает переполнение памяти при работе с «тяжелыми» JSON-документами.
Kafka & Databus
Автоматическая сериализация в Protobuf с проверкой схем. Для Kafka реализован контроль переполнения буфера продюсера и гарантированная доставка (Ack).
Data Pilot: Математика предсказания (EMA)
Жесткое расписание (Cron) слепо. Оно не знает, есть ли данные в источнике. Data Pilot заменяет расписание на экспоненциальное скользящее среднее (EMA).
Система делит сутки на слоты и для каждого запоминает объем данных. Благодаря свойству «забывания» (коэффициент α), модель быстро адаптируется к новым нормам трафика, игнорируя разовые всплески.
// α — коэффициент забывания, настраиваемый AutoTuner
Математическая фильтрация аномалий (3σ)
Как отличить реальный сбой от временного всплеска активности? Мы используем правило трех сигм. Если фактический объем данных выходит за пределы коридора Mean ± 3σ, система фиксирует аномалию.
Это позволяет Пилоту мгновенно переходить в режим «Турбулентности» или «Форсажа», учащая проверки и увеличивая размер батчей, чтобы не допустить накопления очереди.
4 Автономных режима работы AutoTuner
Пайплайн — это живой организм. В зависимости от статистики он самостоятельно выбирает одну из четырех фаз:
- Крейсерский: Поток стабилен. Оптимальные интервалы, минимальная нагрузка.
- Затишье: Данных меньше нормы. Пилот увеличивает паузы, экономя ресурсы Source и не «молотя» базу впустую.
- Турбулентность (3σ): Резкий всплеск. Система переходит на работу в реальном времени, разгружая очередь.
- Форсаж: После длительного простоя отключает аналитику, чтобы не «отравлять» статистику выбросами, и качает на пределе возможностей.
Изоляция сбоев и отказоустойчивость
Мы отказались от принципа «все или ничего». Если один из приемников (например, OpenSearch) недоступен, фреймворк не останавливает весь поток. Данные успешно уходят в ClickHouse и Kafka.
«Хвост» (диапазон инкремента) для неудавшегося транспорта сохраняется в таблице logmodel.etl_failed_batches. При следующем запуске система видит этот «долг» и выполняет досылку только для проблемного участка, не дергая источник заново.
Гарантия доставки
100% консистентность данных у всех потребителей благодаря механизму Retry Logic без участия человека.
Сводные показатели эффективности
Внедрение фреймворка превратило хаос ручных DAG-ов в управляемую экосистему. Цифры говорят сами за себя:
| Метрика | До внедрения | После внедрения |
|---|---|---|
| Нагрузка на Source БД | Постоянная высокая (холостые опросы) | Снижена на 20–30% в периоды затишья |
| Реакция на пики | Задержка до следующего часа | Мгновенная (в пределах минуты) |
| Поддержка новых DAG-ов | Часы разработки и тестов | Минуты (конфигурация через Dataclass) |
| Ложные алерты | Высокий уровень шума | Подавлены на 95% благодаря 3σ |
| Целостность данных | Риск потерь при сбоях сети | 100% гарантия доставки (Retry Logic) |