data-infra
Словарь ↗Конвейер данных (Data Pipeline)
Конвейер данных (data pipeline) — это набор автоматизированных, как правило последовательных, шагов обработки, которые перемещают данные из одной или нескольких исходных систем (баз данных, API, потоков событий, файловых загрузок) в место назначения (хранилище данных, векторное хранилище, аналитическую панель, другое приложение), применяя по пути трансформации — очистку некорректных записей, изменение формата, обогащение строк дополнительными данными или удаление дубликатов. Почему это важно для разработчиков AI/SaaS: почти каждая AI-функция опирается на конвейер за кулисами, даже когда поверхность продукта выглядит простой. Функция «чат с вашими данными» нуждается в конвейере, который забирает документы оттуда, где они хранятся, разбивает их на чанки и встраивает (embedding), а также поддерживает синхронность векторного хранилища при изменении исходных документов. Функция биллинга по использованию нуждается в конвейере, агрегирующем сырые события в дневные/месячные сводки. Построить это как разовые скрипты работает для демо; построить как наблюдаемые, безопасно повторяемые (retryable), идемпотентные конвейеры — вот что отличает хрупкий побочный проект от продукта, которому клиенты доверяют свои данные. Как это работает: конвейеры обычно описывают как пакетные (batch — обработка накопленных данных по расписанию: ночью, ежечасно) или потоковые (streaming — непрерывная обработка событий по мере поступления, часто через очередь сообщений вроде Kafka или Redis Streams). Большинство промышленных конвейеров оркестрируются планировщиком/DAG-инструментом (Airflow, Dagster, Prefect или более простые cron-запускальщики для небольших команд), который отслеживает зависимости между шагами, повторяет неудачные попытки и оповещает об ошибках вместо тихого падения. Хорошо спроектированный конвейер идемпотентен — повторный запуск на тех же входных данных даёт тот же результат, а не дублирует данные, — что критически важно, когда шаг падает на середине и его нужно безопасно повторить. Наблюдаемость (логирование числа строк на входе и выходе каждого этапа, отслеживание длительности выполнения, оповещение об аномалиях) превращает конвейер из чёрного ящика в то, что команда реально может отладить в два часа ночи. Практический пример: B2B SaaS синхронизирует данные CRM клиента каждую ночь, чтобы формировать AI-сводку о состоянии аккаунта. Конвейер: (1) извлечение — забирает новые/обновлённые записи из API Salesforce с момента последнего успешного запуска, используя сохранённую отметку `last_synced_at`; (2) трансформация — нормализует имена полей, удаляет поля с персональными данными, не нужные далее по цепочке, помечает записи с отсутствующими обязательными полями для очереди недоставленных сообщений (dead-letter queue) вместо их тихого отбрасывания; (3) загрузка — выполняет upsert во внутреннюю таблицу хранилища Postgres и повторно встраивает изменившиеся поля заметок по аккаунту в векторное хранилище. Каждый этап логирует число строк в панель мониторинга, и срабатывает оповещение в Slack, если число извлечённых строк падает более чем на 50% относительно среднего за 7 дней — это позволяет обнаружить сломанный токен API Salesforce до того, как синхронизация данных незаметно остановится на неделю.
Похожие термины