Word Count в YTsaurus Flow (Go)
Простейший пример stateful-пайплайна на Go: подсчёт количества вхождений каждого слова во внутреннем YSON-стейте.
Структура
Пайплайн состоит из двух компьютейшенов:
reader— нативный сорс (TSwiftPassthroughOrderedSourceComputation), объявленный прямо в спеке: он читает очередь и публикует строки в стримwords. Go-кода у него нет.mapper— transform-компьютейшен (TTransformCompanionComputation), который обслуживает компаньон: он читает стримwordsи обновляет счётчик слова во внутреннем стейте.
Сообщения группируются по слову (group_by_schema с farm_hash(word) и word), поэтому стейт обрабатываемого ключа — это счётчик ровно одного слова. Результат пайплайна лежит в таблице внутреннего стейта: дальше по графу ничего не отправляется.
main.go
Точка входа: создание пайплайна, регистрация единственного компьютейшена и запуск.
word_count_mapper.go
Значение стейта — обычная Go-структура с YSON-тегами: именно в этом виде она лежит в таблице внутреннего стейта.
flow.RowFunction, которая открывает внутренний стейт по имени word-state через flow.OpenYSONState и увеличивает счётчик:
Ключевые паттерны
- Простейший stateful-пайплайн с одним компьютейшеном на стороне компаньона: сорс остаётся нативным, и Go-кода для него не требуется.
- Внутренний YSON-стейт через
flow.OpenYSONState[T](rt, name, msg):Value()возвращает изменяемую структуру, которую SDK сохраняет после успешного батча. - Имя стейта (
word-state) совпадает с именем изparameters.internal_statesкомпьютейшена в спеке. - Ключ стейта определяется
group_by_schemaиз спеки — в данном случае по полюword. - Вход один раз преобразуется в
wordMessageчерезmsg.ConvertTo(&input), после чего обработчик работает с полями структуры.