Process function in YTsaurus Flow (C++)

Why you need it

Implement new C++ user logic only as a process function. Direct inheritance from TTransformComputation, TSwiftMapComputation, TSwiftOrderedSourceComputation, or TTransformOrderedSourceComputation is a low-level API for framework and legacy maintenance, not an authoring model for new user computations. Such logic is tightly coupled with the Computation object, inherits dozens of protected methods, and can only be constructed from a fully built TComputationContext, so it’s nearly impossible to test in isolation with unit tests.

A process function moves your custom logic into a separate, lightweight object. This object receives its dependencies (IOutputCollector, IRuntimeContext) as narrow interfaces and doesn’t depend on the Computation object itself. This lets you test the function in isolation with unit tests.

How it works

The process function isn’t run directly. Instead, a built-in Computation adapter executes it. You specify this adapter in the spec via the computation_class_name field (see Registration). The adapter gives the function the same environment as a regular Computation: the same spec, the same stores and states, and the same epoch-processing logic. That’s why a pipeline with a process function provides the same processing guarantees (including exactly-once) as a manually written Computation.

The adapter also defines the function’s execution mode. There are four built-in adapters:

computation_class_name Mode
NYT::NFlow::TProcessFunctionComputation transform
NYT::NFlow::TProcessFunctionSwiftMapComputation swift-map
NYT::NFlow::TProcessFunctionSourceComputation source
NYT::NFlow::TProcessFunctionTransformOrderedSourceComputation ordered source with output materialization and state

You can run the same function under different adapters without rebuilding the binary.

Interfaces

The library is library/cpp/common (common/process_function.h). The function selects one processing granularity by inheriting the corresponding interface. The spec determines which Computation (source, swift map, or transform) the function attaches to (see Registration). The base IProcessFunctionBase only includes Init(initContext) — initialization at the start of an epoch. By default, this is a no-op. The granularity interfaces add the actual processing methods.

Choose the interface based on how you want to process the epoch’s input: one entity at a time, the entire batch at once, or by key. Then override only the methods you need. The worker will call them with the state and exactly-once semantics of the selected mode.

The function inherits one granularity interface and, if needed, the ISyncProcessFunction mix-in:

  • IProcessFunction — element-wise processing (the most common case). The worker calls a method for each entity in the epoch. Override the methods you need; all are no-op by default:

    • ProcessMessage(message, output, context) — handles a single message.
    • ProcessTimer(timer, output, context) — handles a single timer.
    • ProcessVisit(visit, output, context) — handles a single visit.

    In source mode, only messages arrive, so ProcessTimer and ProcessVisit aren’t called.

  • IBatchProcessFunction — processes the entire epoch input in a single call. Override Process(input, output, context) when your logic works with the whole batch at once (for example, a single batched external request). The input isn’t grouped by key. To combine batch work with per-entity handling — say, one state preload for the whole batch — call the dispatch helpers ProcessMessages / ProcessTimers / ProcessVisits(input, output, context, callback) from Process: each sets the parents and tags errors with the key exactly as the worker does around IProcessFunction hooks, and takes a TCallback — BIND(&TMyFunction::ProcessMessage, MakeStrong(this)) or a BIND of a lambda.

  • IKeyedBatchProcessFunction — processes by key using group-by, for keyed modes (swift map and transform). The worker groups the epoch’s input by key and calls ProcessKey for each key:

    • ProcessKey(input, output, context) — handles all input for a single key (messages, timers, and visits together). It’s no-op by default. Override it when your logic relies on the entire key batch at once (for example, to reconcile messages and timers via the key’s shared state).
  • ISyncProcessFunction — an optional mix-in for functions that commit side effects in a separate sync phase at the end of the epoch. You inherit it in addition to the granularity interface:

    • Sync(transaction, context) — commits side effects in the transaction. The context gives access to runtime accessors. You must implement this method. It’s called only by a Computation adapter that has a sync phase — among the built-in adapters, that’s TProcessFunctionComputation (transform) and TProcessFunctionTransformOrderedSourceComputation (ordered source). The spec validation checks this match: you can’t attach a function with Sync to a Computation without a sync phase.

In Process (IBatchProcessFunction) and ProcessKey, the output doesn’t have parent messages set — you must set them yourself via output->SetParents(...). In the element-wise methods of IProcessFunction (ProcessMessage, ProcessTimer, ProcessVisit), they’re already set for the corresponding entity.

The Distribute field in TAddMessageOptions mirrors the OutputCollector semantics in Computation: for source, a message with Distribute = false isn’t published but is still considered when evaluating the watermark; for other Computation types, such a message is simply discarded.

Message ID suffixes in Swift

In Swift mode, you can select the derived message ID suffix in TAddMessageOptions:

output->AddMessage(std::move(message));
output->AddMessage(
    std::move(hashedMessage),
    TAddMessageOptions{
        .MessageIdSuffix = TOutputMessageIdSuffix::FromPayloadHash(),
    });
output->AddMessage(
    std::move(keyedMessage),
    TAddMessageOptions{
        .Distribute = false,
        .MessageIdSuffix = TOutputMessageIdSuffix::FromUserDefined(semanticKey),
    });
  • With default options, or with FromSequenceNumber(), Flow uses the current message sequence number for the parent message ID and output stream pair.
  • FromPayloadHash() uses a 128-bit CityHash of the canonical payload wire representation. Equal payloads with the same parent and stream get the same message ID; a hash collision is theoretically possible.
  • FromUserDefined(...) accepts a non-empty user-defined suffix. Flow encodes it as a separate lexicographic component, so it cannot impersonate a sequence-number suffix. The caller is responsible for keeping it stable and unique among distinct logical messages with the same parent and stream.

Payload hashes and user-defined suffixes are supported only by Swift computations. They let message identity be independent of emission order, but they don’t remove the Swift determinism requirement: the same message ID must still denote the same logical output.

Warning

A process function is TRefCounted, so you must always create it via New<...>().

IRuntimeContext

context (of type IRuntimeContext, common/runtime_context.h) is an interface that collects everything a Computation usually reads from this:

Method Description
GetWatermark(streamId) / GetInputEventWatermark() Event-time watermarks
GetSpec() / GetStreamSpecs() / GetKeySchema() Spec and stream schemas
MakeOutputMessageBuilder(streamId) Output message builder
ConvertToOutputMessage(message, streamId) Converts a message to the output stream’s schema
ConvertToMessage(ysonMessage) TYsonMessage → TMessage
ConvertToYsonMessage<T>(message) TMessage → typed TYsonMessage
MakeTimer(key, streamId, trigger, event) Creates a timer
GetThrottlerOrThrow(throttlerId) Gets a distributed throttler
TryGetThrottler(throttlerId) Same, but returns nullptr if the throttler is not declared

States

Typed state clients (TMutableStateKeyClient<T> and others) are stored as function fields and initialized in Init via IRuntimeInitContext (common/runtime_init_context.h):

void Init(const IRuntimeInitContextPtr& initContext) override
{
    initContext->InitExternalStateClient(StateClient_, "/state");
    // or: initContext->InitClient<TMyState>(Client_, "my-state");
}

For more details on state types, see the Working with states section.

Parameters

A function can declare its own parameter structure — a regular TYsonStruct — and read it in a typed way. You pass parameters in the processing_function_parameters field at the top level of the Computation spec (next to processing_function):

  • Static — processing_function_parameters in spec, read once in Init via initContext->GetParameters<T>().
  • Dynamic — processing_function_parameters in dynamic_spec, read via context->GetDynamicParameters<T>() and reflect the latest reconfiguration.

Both blocks are parsed into the parameter types declared when registering the function (see below), so T in GetParameters<T>() / GetDynamicParameters<T>() must be the registered type — a mismatch throws. If the processing_function_parameters field is missing, the structure is filled with default values. The static block is parsed once at job init; the dynamic one is reparsed only when it changes (that is, on reconfiguration).

You specify parameter types when registering the function via macro arguments: YT_FLOW_DEFINE_PROCESS_FUNCTION(function, TStaticParams) for the static block, YT_FLOW_DEFINE_PROCESS_FUNCTION(function, TStaticParams, TDynamicParams) also for the dynamic one. Then the corresponding processing_function_parameters block (in spec and dynamic_spec) is validated against the schema when loading the spec, just like parameters for Computation: an unknown field or incorrect type causes an error before the run. A block for which you didn’t declare a type (including for a parameterless YT_FLOW_DEFINE_PROCESS_FUNCTION(function)) is treated as empty — any passed field will be rejected.

struct TMyParameters
    : public NYTree::TYsonStruct
{
    i64 Threshold;

    REGISTER_YSON_STRUCT(TMyParameters);

    static void Register(TRegistrar registrar)
    {
        registrar.Parameter("threshold", &TThis::Threshold).Default(0);
    }
};

void Init(const IRuntimeInitContextPtr& initContext) override
{
    Threshold_ = initContext->GetParameters<TMyParameters>()->Threshold;
}

void ProcessMessage(const TInputMessageConstPtr& message, const IOutputCollectorPtr& output, const IRuntimeContextPtr& context) override
{
    auto currentThreshold = context->GetDynamicParameters<TMyParameters>()->Threshold;
    // ...
}

// Registering with a parameter type enables validation in the spec.
YT_FLOW_DEFINE_PROCESS_FUNCTION(TMyFunction, TMyParameters);

In the spec:

"counter" = {
    "computation_class_name" = "NYT::NFlow::TProcessFunctionComputation";
    "processing_function" = "NYT::NFlow::NExample::TWordCountFunction";
    "processing_function_parameters" = {
        "threshold" = 5;
    };
};

In unit tests, set static parameters via TTestStateEnvironment::SetStaticParameters(...) — the passed struct is served to GetParameters<T>() as is. Dynamic ones go through the production path: name the function via TTestRuntimeContextBuilder().SetProcessingFunction<TMyFunction>() (it must be registered in the test binary) and pass the struct to SetDynamicParameters(...); the context reparses it into the registered dynamic type.

Registration

You link the function and Computation via the spec. Register the function with a single macro from common/registry.h (in the same TRegistry as computation / source / sink) under its TypeName; the optional second argument is the type of its parameters (see Parameters):

YT_FLOW_DEFINE_PROCESS_FUNCTION(TWordCountFunction);                   // no parameters
YT_FLOW_DEFINE_PROCESS_FUNCTION(TTextReadFunction, TTextReaderParameters);  // with parameters

In the Computation spec, the computation_class_name field points to the built-in Computation adapter (which also sets the mode — see the list of adapters), and the adjacent processing_function field names the function:

"counter" = {
    "computation_class_name" = "NYT::NFlow::TProcessFunctionComputation";
    "processing_function" = "NYT::NFlow::NExample::TWordCountFunction";
};

Testing

The library/cpp/process_function/testing library provides a ready-made set of utilities for unit tests with sensible defaults (all its utilities live in the separate NYT::NFlow::NTesting namespace to keep test code separate from production):

  • TRecordingOutputCollector — an IOutputCollector that records messages and timers in vectors (GetMessages() / GetTimers()).
  • TTestRuntimeContextBuilder — builds an IRuntimeContext; by default, it uses zero watermarks, one output stream for each registered RegisterStream<T>(id), and the key schema DefaultTestKeySchema().
  • TTestStateEnvironment — starts a TJobStateManager over in-memory mock tables and provides an IRuntimeInitContext; PreloadKeyStates(inputContext) loads keys before processing, and ReadKeyState<T>(name, key) reads the state afterward.
  • entity_builders.h — MakeTestMessage, MakeTestRawMessage, MakeTestTimer, MakeTestVisit.
  • TProcessFunctionTestHarness — runs the function across epochs like a worker does: it wraps the function as batch (via WrapAsBatch, so per-element, whole-batch, and per-key forms all work the same way, and timers and visits are dispatched alongside messages), calls Init once, and for each RunEpoch(...) it preloads the state, processes the input, executes the end-of-epoch Sync for ISyncProcessFunction, and commits the state. You can access the last epoch’s messages and timers via GetMessages() / GetTimers().

The umbrella header process_function/testing/unittest.h includes all the listed utilities; in a test, you only need to include it instead of the individual harness headers.

TTestStateEnvironment covers all state types the function uses:

  • Internal (InitClient) — work directly via TJobStateManager over in-memory tables; the function writes them and reads them back via ReadKeyState<T>(name, key).
  • External managers (InitExternalStateClient with TMutableStateKeyClient<T>) — register via RegisterExternalState(name, ...); there’s a ready-made in-memory TInMemorySimpleExternalStateManager for TSimpleExternalState, and you can read the result via ReadExternalKeyState<T>(name, key).
  • External joiners (InitExternalStateClient with TJoinedStateKeyClient<T>) — register via RegisterExternalStateJoiner(name, ...); the in-memory TInMemorySimpleExternalStateJoiner is seeded via GetMutableState(key) before the function starts.
  • Internal joiners of another computation’s state (InitClient with TJoinedStateKeyClient<T>) — register via RegisterStateJoiner(name, stateName); in a single environment, you can run a producer function (it writes the state), call Sync(), and then run a joiner function that reads it.

RegisterExternalState / RegisterExternalStateJoiner also accept an arbitrary IExternalStateManagerPtr / IExternalStateJoinerPtr — this way, you connect a real manager built over a YTsaurus mock client in the test.

Example of a test for a row function with state:

TTestStateEnvironment stateEnv;
auto context = TTestRuntimeContextBuilder().Build();
auto output = New<TRecordingOutputCollector>();

auto function = New<TCountingRowFunction>();
function->Init(stateEnv.GetInitContext());

auto key = MakeKey<ui64>(7);
auto message = MakeTestMessage("input", key, New<NTableClient::TTableSchema>());
stateEnv.PreloadKeyStates(New<TInputContext>(
    std::vector<TInputMessageConstPtr>{message},
    std::vector<TInputTimerConstPtr>{}));

function->ProcessMessage(message, output, context);

EXPECT_EQ(stateEnv.ReadKeyState<i64>("counter", key), 1);

The same test using TProcessFunctionTestHarness, which hides Init, preload, context/output setup, and epoch commit:

TTestStateEnvironment stateEnv;
TProcessFunctionTestHarness harness(stateEnv, New<TCountingRowFunction>());

auto key = MakeKey<ui64>(7);
harness.RunEpoch({MakeTestMessage("input", key, New<NTableClient::TTableSchema>())});

EXPECT_EQ(stateEnv.ReadKeyState<i64>("counter", key), 1);

For more details on the general approach to testing C++ computations, see the Testing section.

Example

Full example — examples/cpp/word_count (both classes are IProcessFunction in the NYT::NFlow::NExample namespace, registered via YT_FLOW_DEFINE_PROCESS_FUNCTION). In pipeline.yson, TTextReadFunction is connected to TProcessFunctionSourceComputation, and TWordCountFunction is connected to TProcessFunctionComputation via processing_function; TTextReadFunction reads its static min_word_length parameter from processing_function_parameters. A more complex example with timers and window-based joining is examples/cpp/wait_click_join.

See also

Previous