Вход на сайт

Просмотр новости

Найдите то, что Вас интересует

Открываем код YTsaurus Flow: как обрабатывать более 100 ГБ/с в реальном времени без потерь и дублей

Дата публикации: 06-10-2026 07:00:18

Если вы когда‑либо строили пайплайны для обработки потоков данных в реальном времени со строгими требованиями, то наверняка знаете, сколько инфраструктурных заморочек в этом деле. Например, обеспечивать относительно низкие end‑to‑end‑задержки в штатном режиме относительно несложно. Но что, если они нужны и в высоких перцентилях, в том числе при сбоях или плановом обслуживании одного из используемых дата‑центров целиком? Как правильно партиционировать поток данных, обрабатывающие его процессы и их долгосрочное состояние, если нагрузка постоянно меняется? Как гарантировать exactly‑once, то есть отсутствие потерь и дублей даже при сбоях оборудования и в краевых случаях? Как понять, обработали ли мы все данные на тот или иной момент? С такими вопросами мы столкнулись при разработке высоконагруженных рекомендательных систем. Для этих задач мы создали YTsaurus Flow — фреймворк потоковой обработки данных с сохранением состояния между событиями и гарантиями exactly‑once по умолчанию. Это совместный проект команд Yandex Infrastructure и Яндекс Рекламы для обработки потоков данных в реальном времени. И сегодня мы открываем его исходный код под лицензией Apache® 2.0.Flow — часть YTsaurus, платформы хранения и обработки данных, которую мы выложили на GitHub в марте 2023 года. Flow использует хранилище, очереди и общий механизм транзакций платформы, чтобы брать на себя управление состоянием и восстановление обработки после сбоев. Разработчик пайплайна при этом сосредоточивается на прикладной логике.В статье разберём, зачем мы разработали собственный движок, как устроены его пайплайны и за счёт чего обеспечивается exactly‑once. А на примере реального сервиса покажем, как переход на Flow помог сократить задержку поставки данных для дообучения моделей с десятка‑другого часов примерно до двух. Читать далее

Основное содержимое страницы с новостью.

194c2b042e22979a41959885ed00ad88.png

Если вы когда‑либо строили пайплайны для обработки потоков данных в реальном времени со строгими требованиями, то наверняка знаете, сколько инфраструктурных заморочек в этом деле. Например, обеспечивать относительно низкие end‑to‑end‑задержки в штатном режиме относительно несложно. Но что, если они нужны и в высоких перцентилях, в том числе при сбоях или плановом обслуживании одного из используемых дата‑центров целиком? Как правильно партиционировать поток данных, обрабатывающие его процессы и их долгосрочное состояние, если нагрузка постоянно меняется? Как гарантировать exactly‑once, то есть отсутствие потерь и дублей даже при сбоях оборудования и в краевых случаях? Как понять, обработали ли мы все данные на тот или иной момент? 

С такими вопросами мы столкнулись при разработке высоконагруженных рекомендательных систем. Для этих задач мы создали YTsaurus Flow — фреймворк потоковой обработки данных с сохранением состояния между событиями и гарантиями exactly‑once по умолчанию. Это совместный проект команд Yandex Infrastructure и Яндекс Рекламы для обработки потоков данных в реальном времени. И сегодня мы открываем его исходный код под лицензией Apache® 2.0.

Flow — часть YTsaurus, платформы хранения и обработки данных, которую мы выложили на GitHub в марте 2023 года. Flow использует хранилище, очереди и общий механизм транзакций платформы, чтобы брать на себя управление состоянием и восстановление обработки после сбоев. Разработчик пайплайна при этом сосредоточивается на прикладной логике.

В статье разберём, зачем мы разработали собственный движок, как устроены его пайплайны и за счёт чего обеспечивается exactly‑once. А на примере реального сервиса покажем, как переход на Flow помог сократить задержку поставки данных для дообучения моделей с десятка‑другого часов примерно до двух.

Зачем нужен ещё один движок для обработки данных в реальном времени

В миру существует довольно много движков для обработки потоков сообщений. Конечно, не больше тысячи, как СУБД для долгосрочного хранения данных, но тоже прилично. Быть частью более широкой экосистемы для них не редкость: в экосистеме Apache Spark™ есть Spark Structured Streaming, а в Apache Kafka® — Kafka Streams. Даже относительно обособленный Apache Flink® часто доступен как управляемый сервис у облачных провайдеров.

Таким образом, если перед компанией встаёт вопрос о выборе такого движка, то логичнее всего рассматривать его в контексте более широкой инфраструктуры для работы с данными. Но в этой статье мы ограничимся потоковой обработкой в реальном времени — разберём, какие задачи решает YTsaurus Flow и как он устроен.

Основной сценарий, под который изначально проектировался YTsaurus Flow, — построение высоконагруженных real‑time рекомендательных систем. Именно он повлиял на ключевые архитектурные решения и особенности YTsaurus Flow.

Состояние по ключу

Пайплайны рекомендательных систем редко выглядят как простая «труба» из одной очереди сообщений в другую. Как правило, они получают на вход поток событий, происходящих на сервисе, и на его основе им приходится поддерживать долгосрочное состояние (стейт) с эмбеддингами и/или статистикой по неким сущностям. Скажем, по баннерам, если мы говорим о рекламных системах, в случае персонализации — по пользователям и так далее. 

При большом потоке входных данных и, как следствие, нагрузке на пайплайн важно работать со стейтом максимально эффективно. В этой области возможна масса оптимизаций, которые реализованы в YTsaurus Flow и позволяют экономить на порядки больше необходимых ресурсов. Подробнее c частью из них можно ознакомиться в записи нашего доклада.

Корректный результат при сбоях

Если пайплайн работает с критически важными для бизнеса данными, от него, скорее всего, потребуется гарантия exactly‑once: сообщения не должны ни теряться, ни учитываться дважды в результатах вычислений.

Поэтому YTsaurus Flow по умолчанию работает в режиме exactly‑once, с гарантией отсутствия дублей и потерь сообщений. Если в конкретной задаче важнее экономить ресурсы, например при подготовке данных для ML, эту гарантию можно явно ослабить: выбрать at‑least‑once, допускающий дубли, или at‑most‑once, допускающий потери. 

Один из ключевых архитектурных подходов YTsaurus Flow для эффективного обеспечения этих гарантий — атомарное изменение и стейта, и метаданных вроде оффсетов очередей в общей транзакции. Это возможно благодаря наличию в экосистеме YTsaurus и своих очередей и таблиц как для транзакционных, так и для аналитических нагрузок — всё под общим встроенным механизмом транзакций.

Доступность

В отличие от бэкендов сервисов, real‑time процессинги не видны напрямую конечным пользователям сервисов, но недоступность или частичные сбои процессингов всё равно косвенно влияют на пользовательский опыт. Как правило, в таких случаях данные обновляются с задержкой, из‑за чего может снижаться качество рекомендаций и других результатов работы сервиса. 

Поэтому для YTsaurus Flow, как и для многих других технологий Яндекса, бесперебойная работа — базовое требование, даже если целый дата‑центр вышел из строя или недоступен из‑за планового обслуживания.

Сложные графы

Чем сложнее бизнес‑логика пайплайна, тем сложнее может быть его граф вычислений. В Яндексе есть пайплайны YTsaurus Flow из сотен узлов — мы называем их компьютейшнами. А для некоторых задач нужны циклы, которые выходят за рамки классического DAG. 

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

Автоматическая адаптация пайплайна под поток данных

На входе real‑time‑пайплайнов бывают десятки и сотни потоков данных, суммарно достигающих десятков гигабайт и миллионов сообщений в секунду, что делает особо острым вопрос эффективности их обработки. 

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

А в периоды низкой real‑time‑нагрузки, например ночью, удобно, что освобождающиеся вычислительные ресурсы могут быть переиспользованы под другие нужды, например для батч‑MapReduce‑операций YTsaurus.

Время события

Бизнес‑логике таких пайплайнов зачастую важно работать в терминах времени, когда событие изначально произошло, а не поступило в обработку. Это важный аспект, который во многих подобных системах реализован. В этой части механизмы YTsaurus Flow во многом созданы под вдохновением от статьи The Dataflow Model: автоматическое отслеживание вотермарков для работы с отстающими данными, поддержка откладывания сообщений через таймеры и так далее.

Нативный исполняемый код

На закуску — ядро YTsaurus Flow разработано на C++. Понятно, что во многих случаях достаточно и JVM, но на масштабе компиляция в нативный исполняемый код — важная экономия ресурсов. Не зря коммерческие сборки и сервисы на основе упоминавшихся выше опенсорсных проектов от Apache Software Foundation нередко включают в себя полностью переписанное альтернативное ядро системы без завязки на JVM. 

В вопросе разработки бизнес‑логики пайплайнов YTsaurus Flow предоставляет свободу выбора: в критичных для производительности кейсах её можно разрабатывать тоже на C++ (что в альтернативных технологиях, как правило, толком не доступно), для прототипов лучше подойдёт Python®, а на спектре между ними доступны Go, Kotlin, Java™. API для разработки YTsaurus‑Flow‑пайплайнов на всех поддерживаемых языках выглядит похожим образом. Помимо этого, простые пайплайны можно описывать целиком на YQL, разработанном в Яндексе диалекте SQL.


В подобных статьях часто встречаются сравнительные бенчмарки. Мы решили не вдаваться сейчас в этот жанр: у систем обработки данных в реальном времени нет общепризнанного стандартного бенчмарка, условного аналога TPC‑C или TPC‑H/DS из мира баз данных. Кастомный бенчмарк всегда можно подобрать так, чтобы в нём выигрывал любой участник, — проверено на практике. 

Так что если вас интересует именно аспект производительности, то лучший способ делать бенчмарки по‑прежнему актуален — конкретно на своих кейсах и остальной низлежащей инфраструктуре. Благо в текущих реалиях это стало требовать на порядок меньше времени и сил благодаря LLM, особенно для опенсорс‑технологий.

Как устроен YTsaurus‑Flow‑пайплайнНа уровне операционной системыe3a6e8526f611c6f4843108a5d675c2f.png

Технически YTsaurus‑Flow‑пайплайн состоит из двух слоёв, каждый в отдельном приложении:

  • Инфраструктурный слой (основное приложение) содержит ядро YTsaurus Flow. Оно координирует работу распределённых компонентов, обеспечивает гарантии фреймворка и управляет состоянием пайплайна. Разработчику пайплайна не нужно дорабатывать его код: поведение приложения задаётся настройками в статической и динамической спецификациях (его конфигах).

  • Прикладной слой (компаньон) — приложение с бизнес‑логикой пайплайна, написанное на любом из поддерживаемых языков программирования. Оно обрабатывает данные, которые получает от основного приложения, возвращает результат обратно и не хранит собственное состояние.

Основное приложение может быть запущено в трёх режимах:

  • Контроллер может быть запущен в нескольких экземплярах. Сначала они выбирают лидера среди друг друга (leader election). Затем лидер координирует распределённую работу пайплайна: планирует, где и что будет обрабатываться, автоматически определяет оптимальное число партиций для всех этапов, перераспределяет нагрузку, перезапускает части вычислений при проблемах и так далее.

  • Воркер выполняет назначенную ему контроллером работу: читает входящие сообщения и отправляет их на обработку в компаньон. Компаньон обычно запущен рядом с воркером как сайдкар, а для C++ бизнес‑логику можно вкомпилировать прямо в воркер.

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

  • Административный клиент проверяет и применяет спецификацию пайплайна, а также позволяет управлять его состоянием. Обычно его вызывают из CI/CD при выпуске новой версии пайплайна. Его взаимодействие с контроллером происходит через общий механизм — YTsaurus RPC proxy.

В плане способов деплоя этих приложений в целом подойдёт любое Linux‑окружение: от железных серверов до корпоративного облака. Главное требование — сетевой доступ до самого кластера YTsaurus, так как именно через него происходит leader election, сохранение стейта обработки, обмен данными через очереди сообщений и так далее. 

Также есть удобная опция запускать YTsaurus‑Flow‑пайплайны прямо внутри YTsaurus‑кластера через механизм Vanilla‑операций. Это аналог классических операций MapReduce, только долгоживущих и без входных таблиц.

Кластеры YTsaurus для батчевой обработки данных через YQL или MapReduce обычно разворачивают в одном дата‑центре для лучшей сетевой связанности между узлами. При обработке в реальном времени важно минимизировать простои и отставание в обработке данных, поэтому для YTsaurus Flow лучше подходит работа в нескольких дата‑центрах. Обычно используют три дата‑центра, чтобы система могла продолжить работу при отказе одного из них.

Важно, чтобы работа в нескольких дата‑центрах была организована для всех компонентов системы. Контроллеры и воркеры YTsaurus Flow должны присутствовать в каждом дата‑центре, а все используемые динамические таблицы YTsaurus должны быть настроены на один из кросс‑дата‑центровых режимов. Также в каждом дата‑центре нужен запас ресурсов, чтобы продолжать обрабатывать данные с требуемой скоростью, даже если треть кластера недоступна.

На уровне обработки данных

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

Во Flow есть несколько десятков типов компьютейшнов. Для одних достаточно указать в спецификации имя и его настройки. Для других придётся также написать бизнес‑логику в компаньоне. Полный список доступных классов компьютейшнов есть в документации, а здесь рассмотрим основные группы:

  • Source — подают данные на вход пайплайна.

  • Transform — обрабатывают данные с доступом ко всей функциональности YTsaurus Flow.

  • Swift — позволяют экономить на материализации промежуточных данных ценой требования детерминированности вычислений и ограниченной функциональности.

  • Sink — отправляют результаты работы пайплайна потребителям.

Пример топологии графа пайплайна в веб-интерфейсе YTsaurus

Пример топологии графа пайплайна в веб‑интерфейсе YTsaurus

Таким образом, основная задача разработчика пайплайна YTsaurus Flow — спроектировать подобную топологию обработки потоков данных с использованием готовых компьютейшнов и затем реализовать бизнес‑логику на предоставляемых API в соответствии с желаемым итоговым результатом.

Exactly‑once в YTsaurus Flow

Exactly‑once — одно из наиболее востребованных свойств при потоковой обработке данных. Если инфраструктурный слой не защищает от дублей и потерь сообщений, эта задача ложится на прикладную логику и её разработчикам приходится постоянно учитывать инфраструктурные вопросы: нужны ли для корректного результата дедупликация или повторные попытки обработки, как их правильно реализовать и так далее.

Чтобы разработчики пайплайнов могли фокусироваться на своей прикладной задаче, в YTsaurus Flow доступны следующие ключевые механизмы «из коробки»:

  • Атомарные изменения. Состояние, прогресс и данные для поддерживаемых синхронных приёмников фиксируются одной транзакцией — каждую такую итерацию мы называем эпохой. Если обработка прерывается, система либо фиксирует весь результат эпохи, либо не фиксирует его вовсе. Помимо обеспечения гарантий корректности такой подход позволяет не создавать лишнюю транзакционную нагрузку на слой хранения YTsaurus. Для асинхронных внешних приёмников используется подход transactional outbox: сначала данные надёжно сохраняются в YTsaurus, затем отправляются вовне с идентификатором для дедупликации и ретраями.

  • Защита от устаревшего воркера. При сбое или изменении нагрузки контроллер может передать единицу работы другому воркеру. После этого старый воркер уже не сможет зафиксировать свой результат и тем самым продублировать работу нового. Технически это обеспечивается механизмом зависимых (prerequisite) транзакций в слое хранения YTsaurus. Одну транзакцию создаёт контроллер, а воркер для коммита своих изменений обязан указать эту транзакцию контроллера как зависимую. То есть для фиксации изменений она тоже должна быть активна. В результате контроллеру достаточно отменить свою исходную транзакцию и запустить эту же работу на другом воркере с новой.

  • Дедупликация входа. Для внутренних потоков данных YTsaurus Flow отслеживает идентификаторы обработанных сообщений, а для входящих источников — позицию чтения. Прогресс продвигается только вместе с успешной фиксацией эпохи.

Важно понимать, что эти механизмы распространяются только на обработку данных непосредственно через YTsaurus Flow. Внутри компьютейшнов может выполняться произвольный код, в том числе с внешними побочными эффектами, например HTTP‑запросами. За корректность и надёжность таких действий отвечает разработчик пайплайна — например, он может обеспечить их идемпотентность.

Гарантии exactly‑once требуют дополнительных вычислительных ресурсов. Так что, если в конкретной задаче важнее их экономить, YTsaurus Flow позволяет явно отказаться от этих гарантий через настройки в спецификации.

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

Разберём, в каких ситуациях YTsaurus Flow может быть полезен, на конкретном примере. 

Чтобы показывать максимально релевантные объявления, рекламной системе важно регулярно дообучать свои рекомендательные ML‑модели на событиях, связанных со свежими манипуляциями взаимодействующих с ними пользователей. Например, увидеть, что пользователю показали баннер в определённом контексте, через какое‑то время он на него кликнул, попал на другой сайт и в результате либо сделал, либо не сделал на нём целевое действие вроде покупки. 

Рекламодатели, их сайты и объявления, пользователи и их интересы, а также другие связанные данные постоянно меняются. Поэтому единожды обученная модель со временем «отстаёт от жизни», и качество рекомендаций начинает падать. Чтобы учитывать эти изменения, модель важно регулярно дообучать.

Для начала давайте рассмотрим упрощённую схему того, как одна из ключевых частей такого процесса была устроена в Яндекс Рекламе до внедрения YTsaurus Flow:

Для начала давайте рассмотрим упрощённую схему того, как одна из ключевых частей такого процесса была устроена в Яндекс Рекламе до внедрения YTsaurus Flow:

867d6cb2a9a05e3ac269a8ccae28f954.png
  1. На входе есть две очереди сообщений:

    • лог самих пользовательских событий;

    • асинхронно посчитанные ML‑факторы по ним.

  2. Инфраструктурный сервис нарезает эти очереди по часам и складывает в статические таблицы YTsaurus, которые заточены на эффективное чтение неизменяемых данных. Получается, что одна таблица содержит час сообщений одной очереди. Этот процесс оптимизирован скорее под долгосрочное хранение данных. При этом промежуточные копии данных несколько раз записываются на диск и считываются с него, что создаёт большие накладные расходы по I/O.

  3. По одной такой таблице событий не понять, случился ли клик для определённого показа объявления: он мог быть на границе часа и попасть в соседние таблицы. Таким образом, для корректной обработки данных за час Х приходится дожидаться и таблицы за час Х + 1, чтобы сделать self‑join по идентификаторам показов через MapReduce. В теории можно было бы дожидаться опаздывающих событий и дольше, но исторически такого требования не было. Дальше из результатов этого объединения нужно отфильтровать события, которые должны использоваться для обучения нужных моделей.

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

  5. Результат всей этой цепочки преобразований передаётся уже непосредственно в дообучение, в подробности которого здесь вдаваться не будем.

End‑to‑end‑задержка такого батч‑пайплайна от возникновения события до его влияния на модель составляла десяток‑другой часов в зависимости от обстоятельств. Помимо ожидания готовности исходных таблиц, значительный вклад в увеличение этого времени вносили тяжёлые джойны.

А теперь рассмотрим, как процесс преобразился при переходе на real‑time с внедрением YTsaurus Flow:

278ef06820cd845ac77f5bf973468d55.png
  1. На входе те же очереди.

  2. Очередь событий сразу читаем пайплайном YTsaurus Flow, делаем фильтрацию и распределённую группировку по тем же ключам, по которым объединяли записи в батч‑схеме, и записываем в долгосрочный стейт потенциально нужные для дообучения данные. За счёт распределённого репартиционирования по ключу (distributed shuffle) клик попадёт в ту же запись в стейте, что и показ, к которому он относится.

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

В результате новая схема на YTsaurus Flow позволила сократить end‑to‑end‑задержку примерно до двух часов, то есть почти на порядок, из которых один час — то самое ожидание клика. Точные цифры, как это повлияло на рекламную систему, не публичны, но в итоге выиграли все: пользователи быстрее начинают видеть более релевантные объявления, а рекламодатели получают более качественных лидов.

Что дальше

YTsaurus Flow берёт на себя инфраструктурную часть потоковой обработки: управление состоянием, распределение нагрузки и восстановление после сбоев. В статье мы показали, как устроены эти механизмы, где действуют гарантии exactly‑once и пример, как Flow помог сократить задержку поставки данных для дообучения рекламных моделей примерно до двух часов.

Чтобы попробовать Flow на практике, начните с документации YTsaurus, а затем переходите к документации YTsaurus Flow. Исходный код и релизы движка доступны в основном репозитории YTsaurus — не стесняйтесь присылать пул‑реквесты.

Если для ваших задач потоковой обработки во Flow чего‑то не хватает, расскажите об этом в комментариях, публичном чате YTsaurus или создайте issue в GitHub. Конкретные сценарии помогут нам определить, какие возможности развивать в первую очередь.

Схожие новости

#Наименование новостиТональностьИнформативностьДата публикации
1Как на ровном месте сэкономить 1000+ ядер, или Куда на самом деле уходили 80% CPU09.1208-09-2026
2Как я писал сервер и нечаянно пробил 1М RPS01031-08-2026
3RuntimeNodes: как мы, ML‑щики, стали писать рантайм07.7424-09-2026
4Как на ровном месте сэкономить 1000+ ядер, или Куда на самом деле уходили 80% CPU09.1208-09-2026
5Миграция без права на ошибку: как перенести 70 кластеров MongoDB в 7 и не сломать продакшен010.6502-09-2026
665 бесплатных уроков октября: от LLM и Kubernetes до микросервисов, Kafka и безопасности011.9501-10-2026
7**От пара к Python, или Как космический шаттл, джедаи и ...05.430-09-2026
8Запустили полнотекстовый и гибридный поиск в YDB: рассказываем, что под капотом06.9302-10-2026
9Запустили полнотекстовый и гибридный поиск в YDB: рассказываем, что под капотом06.9302-10-2026
10What is a charm?08.0623-09-2026

Классификация: . Схожих патентов: 0. Схожих новостей: 10. Тональность: 0. Информативность: 7.54. Источник: habr.com.