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
ProcessTimerandProcessVisitaren’t called. -
IBatchProcessFunction— processes the entire epoch input in a single call. OverrideProcess(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 helpersProcessMessages/ProcessTimers/ProcessVisits(input, output, context, callback)fromProcess: each sets the parents and tags errors with the key exactly as the worker does aroundIProcessFunctionhooks, and takes aTCallback—BIND(&TMyFunction::ProcessMessage, MakeStrong(this))or aBINDof 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 callsProcessKeyfor 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 thetransaction. Thecontextgives access to runtime accessors. You must implement this method. It’s called only by aComputationadapter that has a sync phase — among the built-in adapters, that’sTProcessFunctionComputation(transform) andTProcessFunctionTransformOrderedSourceComputation(ordered source). The spec validation checks this match: you can’t attach a function withSyncto aComputationwithout 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_parametersinspec, read once inInitviainitContext->GetParameters<T>(). - Dynamic —
processing_function_parametersindynamic_spec, read viacontext->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— anIOutputCollectorthat records messages and timers in vectors (GetMessages()/GetTimers()).TTestRuntimeContextBuilder— builds anIRuntimeContext; by default, it uses zero watermarks, one output stream for each registeredRegisterStream<T>(id), and the key schemaDefaultTestKeySchema().TTestStateEnvironment— starts aTJobStateManagerover in-memory mock tables and provides anIRuntimeInitContext;PreloadKeyStates(inputContext)loads keys before processing, andReadKeyState<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 (viaWrapAsBatch, so per-element, whole-batch, and per-key forms all work the same way, and timers and visits are dispatched alongside messages), callsInitonce, and for eachRunEpoch(...)it preloads the state, processes the input, executes the end-of-epochSyncforISyncProcessFunction, and commits the state. You can access the last epoch’s messages and timers viaGetMessages()/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 viaTJobStateManagerover in-memory tables; the function writes them and reads them back viaReadKeyState<T>(name, key). - External managers (
InitExternalStateClientwithTMutableStateKeyClient<T>) — register viaRegisterExternalState(name, ...); there’s a ready-made in-memoryTInMemorySimpleExternalStateManagerforTSimpleExternalState, and you can read the result viaReadExternalKeyState<T>(name, key). - External joiners (
InitExternalStateClientwithTJoinedStateKeyClient<T>) — register viaRegisterExternalStateJoiner(name, ...); the in-memoryTInMemorySimpleExternalStateJoineris seeded viaGetMutableState(key)before the function starts. - Internal joiners of another computation’s state (
InitClientwithTJoinedStateKeyClient<T>) — register viaRegisterStateJoiner(name, stateName); in a single environment, you can run a producer function (it writes the state), callSync(), 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.