Яндекс открыл исходный код YTsaurus Flow — фреймворка для потоковой обработки данныхОн обрабатывает события по мере поступления, а не ждёт, пока они накопятся для пакетной обработки. Инфраструктурную часть фреймворк берёт на себя: долгосрочное состояние хранит в динамических таблицах YTsaurus, а число партиций для каждого этапа подбирает автоматически и сам перераспределяет нагрузку между машинами. Код
открыт под Apache 2.0.
Как это работает на практике, видно по кейсу Рекламы. Чтобы дообучать модели, нужно понять, кликнули ли по показанному объявлению, и приложить к этому ML-факторы. Раньше события раскладывали по почасовым таблицам, а клик мог попасть в соседний час, поэтому self-join показов и кликов через MapReduce делали, только дождавшись следующей таблицы. Уже при нарезке на почасовые таблицы промежуточные копии данных по несколько раз писались на диск и читались обратно.
Теперь события группируются по ключу на лету, и клик попадает в ту же запись состояния, что и показ. Тяжёлую очередь факторов (сотни терабайт в час в сжатом виде) в состояние не пишут: её читают с отставанием на час, когда клик, если он был, уже учтён.
В итоге задержка сократилась с десятка-другого часов примерно до двух, почти на порядок, и без лишней записи огромных объёмов данных.
Базы данных