Shuffle в YTsaurus Flow (C++)
Пайплайн читает поток в формате JSON из сортированной динамической таблицы, многократно группирует по разным ключам, а после считает число уникальных ключей во всех получившихся потоках. Для пайплайна также описан тест.
Пайплайн не пытается решить какую-либо бизнесовую задачу.
Общее описание пайплайна
Reader. Чтение входных данных
Первый компьютейшн - Reader. Он занимается чтением и первичным преобразованием данных из очереди. Вот части спеки, связанные с данным компьютейшном:
{
"spec" = {
"computations" = {
"reader" = {
"computation_class_name" = "NYT::NFlow::TProcessFunctionSourceComputation";
"processing_function" = "NYT::NFlow::NExample::TQueueReader";
"output_stream_ids" = ["event"];
"source_streams" = {
"queue" = {
"source_class_name" = "NYT::NFlow::TQueueSource";
"parameters" = {
"queue_path" = "<cluster=cluster_name>//path/to/queue";
"consumer_path" = "<cluster=cluster_name>//path/to/consumer";
"finite" = false;
};
};
};
};
};
"streams" = {
"event" = {
"schema" = [
{"name" = "key_a"; "type" = "uint64";};
{"name" = "key_b"; "type" = "uint64";};
{"name" = "key_c"; "type" = "uint64";};
{"name" = "key_d"; "type" = "uint64";};
{"name" = "value"; "type" = "string";};
];
}
};
};
}
Разберем детально.
computations/reader/source_streamsсодержит источникqueueс типомNYT::NFlow::TQueueSource. Данныйsourceпредназначен для чтения данных из сортированной динамической таблицы с использованиемconsumer. Вparametersуказывается из какой очереди и каким консьюмером необходимо читать данные. С параметрами детальнее можно познакомиться в рамках классаNYT::NFlow::TQueueSourceParameters.- Для управления
Computationнеобходимо использоватьNYT::NFlow::TQueueSourceController- так как нам нужно определять число и настройки партиций на базе входной сортированной динамической таблицы. streamsсодержит один потокevent- распаршенный поток на выходе изreader, доступный другимComputation. Для него описана соответствующая схема. Этот же поток зарегистрирован и вcomputations/reader/output_stream_ids.- Класс
TQueueReaderреализуетIProcessFunction; его запускаетTProcessFunctionSourceComputation, указанный вcomputation_class_name. Встроенный passthrough здесь не подходит, поскольку для разбораJSONнужна пользовательская реализацияProcessMessage.
- Source-адаптер работает в Swift-режиме: выходные потоки не материализуются в YTsaurus, сохраняется только метаинформация, необходимая для детерминированной работы.
- Доступом к очереди, включая нелокальный кластер, управляют
TQueueSourceи адаптер; process function получает уже прочитанное сообщение и не обращается к клиенту YTsaurus напрямую.
Shuffle
В пайплайне присутствует несколько перемешиваний: shuffle_a, shuffle_b, shuffle_c, shuffle_d. Каждый из них группирует входной поток по соответствующему ключу - key_a, key_b, key_c или key_d. Никаких преобразований с данными они не делают, лишь демонстрируют возможность сгруппировать разные объекты.
Разберем спеку на примереshuffle_b:
{
"spec" = {
"computations" = {
"shuffle_b" = {
"computation_class_name" = "NYT::NFlow::TSwiftPassthroughComputation";
"group_by_schema" = [
{"name" = "hash"; "expression" = "farm_hash(key_b)"; "type" = "uint64";};
{"name" = "key_b"; "type" = "uint64";};
];
"input_stream_ids" = ["event_a"];
"output_stream_ids" = ["event_b"];
};
};
streams = {
"event_a" = {
"schema" = [
{"name" = "key_a"; "type" = "uint64";};
{"name" = "key_b"; "type" = "uint64";};
{"name" = "key_c"; "type" = "uint64";};
{"name" = "key_d"; "type" = "uint64";};
{"name" = "value"; "type" = "string";};
];
};
"event_b" = {
"schema" = [
{"name" = "key_a"; "type" = "uint64";};
{"name" = "key_b"; "type" = "uint64";};
{"name" = "key_c"; "type" = "uint64";};
{"name" = "key_d"; "type" = "uint64";};
{"name" = "value"; "type" = "string";};
];
};
};
};
}
- Так как в рамках примера нет какого-либо преобразования данных, то нам достаточно
NYT::NFlow::TSwiftPassthroughComputation. Если преобразование понадобится, пользовательскую логику следует реализовать какIProcessFunctionи запустить черезNYT::NFlow::TProcessFunctionSwiftMapComputation. NYT::NFlow::TSwiftPassthroughComputationне материализует данные в YTsaurus.group_by_schemaсодержит соответствующий ключkey_b. В него добавлена колонкаhash, так как партиционирование воFlowработает только в предположении, что первая колонка содержит равномерно распределенные значения типаuint64.input_stream_idsиoutput_stream_idsсодержат соответственноevent_aиevent_b.- В
spec/streamsтакже содержатсяevent_aиevent_bс описанием схемы целиком.
Reduce
Последний Computation. Он читает потоки event_a, event_b, event_c, event_d и подсчитывает число встреч каждого value. Фактически, исходный поток обрабатывается четыре раза.
{
"spec" = {
"computations" = {
"reducer" = {
"computation_class_name" = "NYT::NFlow::TProcessFunctionComputation";
"processing_function" = "NYT::NFlow::NExample::TReducer";
"group_by_schema" = [
{"name" = "hash"; "expression" = "farm_hash(value)"; "type" = "uint64";};
{"name" = "value"; "type" = "string";};
];
"input_stream_ids" = ["event_a"; "event_b"; "event_c"; "event_d";];
"output_stream_ids" = [];
"external_state_managers" = {
"/state" = {
"external_state_manager_class_name" = "NYT::NFlow::TSimpleExternalStateManager";
"parameters" = {
"path" = "//path/to/state";
};
};
};
};
};
};
};
- Для описания логики используется process function
TReducer, запущенная черезTProcessFunctionComputation. - Для работы со стейтом используется
TSimpleExternalStateManager, который предоставляет прямой доступ к таблице.TReducerхранитTMutableStateKeyClient<TSimpleExternalState>и привязывает его к"/state"в методеInit(const IRuntimeInitContextPtr&).
См. также
-
Для группировки по
valueобязательно необходимо указать эту колонку (и хэш от неё) вgroup_by_schema. -
В
input_stream_idsперечисляются все потоки:event_a,event_b,event_c,event_d— чтобы читать все получившиеся потоки. С точки зрения "бизнес логики" это не самое осмысленное действие, однако исходной целью данного пайплайна было протестировать гарантииexactly-onceдаже в случаеSwiftцепочки. -
TReducerреализуетIProcessFunctionи запускается черезTProcessFunctionComputation. Адаптер сохраняетinput_messagesиoutput_messagesв YTsaurus, но в этом примере выходных сообщений нет. По сути, пайплайн сохраняет метаинформациюreader, метаинформацию (message_idиkey) каждого входного сообщенияreducerи таблицуvalue => count. Промежуточные passthrough-компьютейшены с YTsaurus не взаимодействуют.
DynamicSpec
- Поле
dynamic_spec/computations/<computation_id>/parameters/desired_partition_countзаполняется для каждогоcomputation, кромеreader. В рамках тестаtest_shuffle.pyпроисходит изменение числа партиций. - В
dynamic_spec/job_tracker/job_threadsуказывается необходимое число тредов для выполнения всех джобов.
Config для запуска
- Ключевое для запуска:
cluster_url,proxy_role,path,rpc_port,monitoring_port. controller/scheduler_periodвыставлен в 200 для конкретного теста - в реальности должно быть достаточно дефолтного значения.logging- настройки логирования.
{
"cluster_url" = "cluster_name";
"path" = "//path/to/pipeline";
"rpc_port" = 81;
"monitoring_port" = 80;
"controller" = {
"scheduler_period" = 200;
};
"logging" = {
"suppressed_messages" = [
];
"rules" = [
{
"exclude_categories" = [
"Bus";
"Dns";
"Concurrency";
"QueryClient";
"Profiling";
"RpcClient";
"Monitoring";
"Net";
"Solomon";
"Jaeger";
"RpcProxyClient";
"RpcServer";
"Dns";
"BufferMetrics";
];
"min_level" = "debug";
"writers" = [
"Stderr";
];
};
];
"writers" = {
"Stderr" = {
"type" = "file";
"file_name" = "/path/to/file.log";
};
};
}
}