Log Parser в YTsaurus Flow (C++)
Пример показывает process function в materialized ordered-source режиме: пайплайн читает строки лога из очереди и парсит их сразу при чтении источника, без промежуточного passthrough-компьютейшена. Функцию исполняет встроенный TProcessFunctionTransformOrderedSourceComputation. Дополнительно пример показывает чтение из Source и собственный durable-стейт, переживающий рестарты (см. Стейт).
Компоненты пайплайна
TLogParserProcessFunction
Пользовательская логика написана как process function — наследник IProcessFunction, который не зависит от объекта Computation, поэтому его можно покрыть юнит-тестами без кластера (unittest). Исполняет её встроенный адаптер TProcessFunctionTransformOrderedSourceComputation, он же задаёт режим — ordered source (см. список адаптеров).
В ProcessMessage(const TInputMessageConstPtr& message, const IOutputCollectorPtr& output, const IRuntimeContextPtr& context) функция читает колонку line сырого сообщения source через GetColumnValue<std::string>(message, "line") и разбирает её с помощью ParseLogLine на записи вида "level:text", разделённые ;, отбрасывая записи без разделителя :, с пустым текстом или с уровнем, отличным от info, warning и error. На каждую валидную запись она обновляет стейт через аксессор StateClient_.GetState(message->Key) (см. Стейт), собирает TLogRecordMessage и эмитит его в выходной стрим records вызовом output->AddMessage(context->ConvertToMessage(outputRecord)).
Результат трансформации — стрим records — материализуется в YTsaurus, как у TTransformComputation, поэтому требований к детерминированности трансформации нет: после рестарта Flow дораспределяет уже материализованные сообщения с ранее назначенными им MessageId, а не вычисляет их заново.
Пример Proto Parser использует тот же режим исполнения через переиспользуемую базу process function: TProtoLogParserFunction наследуется от TProtoParsingProcessFunctionBase<TLogRecordProto> и запускается под TProcessFunctionTransformOrderedSourceComputation. Валидатор спеки у адаптера-хоста тот же, что у базового класса компьютейшена: непустой group_by_schema, таймеры, key-visitor-стримы и external_state_managers отвергаются (полный список ограничений).
Спека компьютейшена parser
Функцию с адаптером связывают два поля спеки компьютейшена parser (см. Регистрация):
"computation_class_name" = "NYT::NFlow::TProcessFunctionTransformOrderedSourceComputation";
"processing_function" = "NYT::NFlow::NExample::TLogParserProcessFunction";
Остальные поля записи parser описывают подключения: source_streams.queue — TQueueSource с путями до очереди и консьюмера, sinks.queue — прямой внешний TSyncQueueSink для стрима records, поэтому отдельного sink-компьютейшена не требуется. Полный файл — pipeline.yson.
Типы сообщений
TLogRecordMessage — наследник TYsonMessage (YSON-структура, регистрируется через YT_FLOW_DEFINE_YSON_MESSAGE) с полями:
level— уровень записи (info,warningилиerror);text— текст записи;worst_level_so_far— максимальный по серьёзности уровень (info < warning < error), встреченный в этой партиции источника на момент записи (см. Стейт).
Стейт
TLogParserProcessFunction — стейтовая. Стейт TWorstSeverityState она держит в поле TMutableStateKeyClient<TWorstSeverityState> StateClient_ (см. Работа со стейтами (C++)). Адаптер TProcessFunctionTransformOrderedSourceComputation вызывает Init(const IRuntimeInitContextPtr& initContext), где клиент подключается к стейту вызовом initContext->InitClient(StateClient_, WorstSeverityStateName) (имя стейта — worst_severity), и ProcessMessage, где стейт читается аксессором GetState(message->Key), а выходные записи приводятся к сообщениям через context->ConvertToMessage(...).
Инстанс компьютейшена привязан к единственной партиции источника, поэтому все сообщения несут один и тот же ключ и обращаются к одной строке стейта: state->WorstSeverity = std::max(state->WorstSeverity, SeverityRank(record.Level)).
Фреймворк синхронизирует этот стейт в той же транзакции эпохи, что и продвижение смещения источника (см. Computation) — сама функция для этого ничего не делает. Поэтому пользовательский стейт корректен exactly-once и не обязан быть идемпотентен к повторной обработке — обычный, «наивный» счётчик обработанных записей был бы здесь так же корректен. Пример хранит именно бегущий максимум серьёзности просто потому, что это естественная агрегатная величина для такого пайплайна, а не потому, что она чем-то безопаснее счётчика.
Функция main
В main выполняется:
NYT::NFlow::Initialize(argc, argv)— инициализация библиотеки Flow.TSimpleSpecBuilder— билдер для регистрации потоков. ЧерезRegisterStream<TLogRecordMessage>("records")регистрируется потокrecordsс типом сообщенийTLogRecordMessage.TSimpleRunnerProgram(std::move(builder)).Run(argc, argv)— запуск пайплайна.
Регистрировать функцию в main не нужно: макрос YT_FLOW_DEFINE_PROCESS_FUNCTION(TLogParserProcessFunction) стоит на файловом уровне в lib/log_parser_process_function.cpp, а сам файл подключён к библиотеке как GLOBAL, поэтому запись в реестре появляется при инициализации бинаря.