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&).

См. также

  • Быстрый старт (C++)

  • Computation (C++)

  • Для группировки по 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";
            };
        };
    }
}
Предыдущая
Следующая