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 выполняется:

  1. NYT::NFlow::Initialize(argc, argv) — инициализация библиотеки Flow.
  2. TSimpleSpecBuilder — билдер для регистрации потоков. Через RegisterStream<TLogRecordMessage>("records") регистрируется поток records с типом сообщений TLogRecordMessage.
  3. TSimpleRunnerProgram(std::move(builder)).Run(argc, argv) — запуск пайплайна.

Регистрировать функцию в main не нужно: макрос YT_FLOW_DEFINE_PROCESS_FUNCTION(TLogParserProcessFunction) стоит на файловом уровне в lib/log_parser_process_function.cpp, а сам файл подключён к библиотеке как GLOBAL, поэтому запись в реестре появляется при инициализации бинаря.

Исходный код

TLogParserProcessFunction

ParseLogLine

См. также

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