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().
См. также
Предыдущая
Следующая