Companion in YTsaurus Flow
In Flow, you can run user code in a separate process. This process is called a companion process.
Companion use cases
Currently used
- Supporting computations in languages other than C++, such as Python, Java and Kotlin, and Go.
- Running user C++ code in a separate process through the C++ companion.
Planned
- Hot updating user code without stopping the pipeline.
Workflow
When you use a companion, the Computation consists of two parts: a specialized Computation on the Worker side and a lightweight Computation on the companion side.
Note
You develop all business logic in the chosen programming language on the companion side, while the pipeline structure is still configured via the spec. In this workflow, the worker becomes an infrastructure binary that doesn’t depend on the pipeline logic. So, when you use Python, Go, Java, or Kotlin, you don’t need to write any C++ user code.
The Computation on the Worker side collects a batch of messages, enriches it with all the information needed for processing (states, parameters, watermark values, and so on), and sends it to the companion via gRPC locally, within a single host.
The batch is formed without regard to keys — one request may contain messages with different keys; this is how the worker collects batches for all computations. There is no per-key grouping in the companion protocol: if the business logic needs per-key processing, it is done in the companion code — see the Python example.
The companion returns its output in groups; each group carries lineage — the list of ids of the input messages of the batch its output was derived from (lineage is unrelated to keys). For when lineage must be set explicitly and what exactly to pass, see When to set lineage explicitly.
In the future, you’ll also be able to use Unix sockets.
You manage the companion process through the resource CompanionManager.
Configuration
Here’s an example of declaring the resource in a static spec for Java:
"CompanionManager" = {
"resource_class_name" = "NYT::NFlow::NCompanion::TJavaCompanionManager";
"parameters" = {
"timeout" = "10s";
"jdk_bin_path" = "/app/ytflow/jdk/bin/java";
"main_class" = "tech.ytsaurus.flow.examples.wordcount.WordCountApplication";
"classpath" = "/app/ytflow/lib/*";
};
"dependencies" = {};
};
For a detailed description of all TCompanionManagerParameters parameters, see the section CompanionManager resource configuration.
The Python companion configuration is described in the section CompanionManager resource configuration (Python).
Here’s an example of declaring a Computation in a static spec:
"computations" = {
"mapper" = {
"computation_class_name" = "NYT::NFlow::NCompanion::TTransformCompanionComputation";
"group_by_schema" = [
{"name" = "hash"; "expression" = "farm_hash(word)"; "type" = "uint64"; required = %true;};
{"name" = "word"; "type" = "string";};
];
"input_stream_ids" = ["words"];
"output_stream_ids" = [];
"required_resource_ids" = {
"CompanionManager" = {
"controller" = false;
"worker" = true;
};
};
"parameters" = {
"internal_states" = ["word-state"];
};
};
};
The key point in this example is using the CompanionManager resource to launch the companion process and the specialized Computation class NYT::NFlow::NCompanion::TTransformCompanionComputation.
C++ companion
You can move user C++ code out of the worker into a separate process as well. The SDK is located in yt/yt/flow/library/cpp/companion/server: you declare the served Computations in TPipeline, specifying the process function type (this typed declaration replaces YT_FLOW_DEFINE_PROCESS_FUNCTION), and build a separate binary with the RunCompanionMain entry point:
int main(int argc, const char** argv)
{
NYT::NFlow::NCompanionServer::TPipeline pipeline;
pipeline.AddSource<TMyReadFunction, TMyReadParameters>("reader");
pipeline.AddTransform<TMyMapFunction>("mapper");
return NYT::NFlow::NCompanionServer::RunCompanionMain(argc, argv, std::move(pipeline));
}
The function is selected by the name from the processing_function field of the Computation spec, the same way as in the in-process TProcessFunctionComputation adapters. The worker launches the binary through the generic TCompanionManager resource:
"CompanionManager" = {
"resource_class_name" = "NYT::NFlow::NCompanion::TCompanionManager";
"parameters" = {
"entrypoint" = {
"executable" = "/path/to/my_companion";
};
};
};
A process function in a companion can take the HTTP clients from IRuntimeInitContext in Init (GetHttpClient(), GetHttpsClient()). The clients are configured by the http_client_config, https_client_config and http_poller_threads fields of the companion block of the worker node config (TCompanionConfig).
Limitations of the first version of the C++ companion:
- sync process functions aren’t supported (the companion protocol has no Sync phase);
- static resources, distributed throttlers, and the epoch timestamp (
GetCurrentTimestamp) aren’t available; GetStreamSpecs()->ComputeKey()can’t compute a key whengroup_by_schemahas computed columns: a companion doesn’t evaluate expressions. The key arrives with the message — usemessage->Key;- external states are supported only as
TSimpleExternalState; - output timers can only reference the key of one of the parent entities of the batch;
- the companion runs as a single multithreaded process (
companion_process_countis 0 or 1).
For an example, see yt/yt/flow/examples/cpp/companion_word_count.
Types of Computations for working with companions
NYT::NFlow::NCompanion::TSwiftMapCompanionComputation: An implementation of TSwiftMapComputation that delegates data processing to the companion process.NYT::NFlow::NCompanion::TSwiftOrderedSourceCompanionComputation: An implementation of TSwiftOrderedSourceComputation that delegates data processing to the companion process.NYT::NFlow::NCompanion::TTransformCompanionComputation: An implementation of TTransformComputation that delegates data processing to the companion process.NYT::NFlow::NCompanion::TTransformOrderedSourceCompanionComputation: An implementation of TTransformOrderedSourceComputation that delegates data processing to the companion process.
Two modes are available for a Source computation. TSwiftOrderedSourceCompanionComputation doesn’t materialize the output and requires deterministic processing without user state. TTransformOrderedSourceCompanionComputation materializes the output and commits it together with the internal state and the source offset in the epoch transaction. Choose it for non-deterministic processing or for working with internal state; the key of such a state is the source partition key. The spec limitations are the same as for TTransformOrderedSourceComputation.
For more details on implementing pipelines using companions, see Java and Kotlin, Python, and Go.