Shuffle в YTsaurus Flow (Python)

Пример пайплайна из двух компьютейшенов: source-компьютейшен парсит JSON-данные и отправляет типизированные сообщения, transform-компьютейшен ведёт подсчёт событий во внешнем стейте.

Исходный код

Структура

  • reader (source) -- EventMapper: парсит JSON из поля data и отправляет типизированные сообщения в стрим event.
  • reducer (transform) -- EventReducer: подсчитывает количество событий по ключу с использованием external state.

__main__.py

{% code '/yt/yt/flow/examples/python/shuffle/main.py' lang='python' lines='[BEGIN main]-[END main]' %}

event_mapper.py

Source-функция, которая парсит JSON из поля data входного сообщения и создаёт типизированное сообщение через ctx.message_builder():

event_reducer.py

Transform-функция с external state для подсчёта событий:

Ключевые паттерны

  • Пайплайн из нескольких компьютейшенов: source + transform.
  • Source-компьютейшен с source=True для чтения из внешнего источника.
  • Парсинг JSON и создание типизированных сообщений через ctx.message_builder().
  • External state с паттерном to_builder() / set() / finish().

См. также

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