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

Точка входа: создание пайплайна, регистрация единственного компьютейшена и запуск.

Исходный код: 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), после чего обработчик работает с полями структуры.

См. также

Предыдущая
Следующая