---
metadata:
  - name: generator
    content: Diplodoc Platform v5.50.6
alternate:
  - https://ytsaurus.tech/docs/en/flow/cpp/process-functions.md
  - https://ytsaurus.tech/docs/ru/flow/cpp/process-functions.md
---
> **Documentation Index:** Fetch the complete configuration index at https://ytsaurus.tech/docs/en/llms.txt

<!-- source: en/_includes/flow/cpp/process-functions.md -->
# 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 {#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](#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](https://ytsaurus.tech/docs/en/flow/concepts/guarantees.md) (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](computation.md#ttransformorderedsourcecomputation) |

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](#registration)). The base `IProcessFunctionBase` only includes `Init(initContext)` — initialization at the start of an [epoch](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#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:

```mermaid
classDiagram
    class TRefCounted {
        <<refcounted>>
    }
    class IProcessFunctionBase {
        <<interface>>
        +Init(initContext)
    }
    class IProcessFunction {
        <<interface>>
        +ProcessMessage(message, output, context)
        +ProcessTimer(timer, output, context)
        +ProcessVisit(visit, output, context)
    }
    class IBatchProcessFunction {
        <<interface>>
        +Process(input, output, context)
    }
    class IKeyedBatchProcessFunction {
        <<interface>>
        +ProcessKey(input, output, context)
    }
    class ISyncProcessFunction {
        <<mix-in>>
        +Sync(transaction, context)
    }
    class TUserFunction {
        <<example>>
    }

    TRefCounted <|-- IProcessFunctionBase
    IProcessFunctionBase <|-- IProcessFunction : element-wise
    IProcessFunctionBase <|-- IBatchProcessFunction : entire epoch
    IProcessFunctionBase <|-- IKeyedBatchProcessFunction : by key
    IProcessFunction <|-- TUserFunction : granularity
    ISyncProcessFunction <|.. TUserFunction : sync 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](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#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](https://ytsaurus.tech/docs/en/flow/cpp/state.md#external-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](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#timestamps-and-watermarks); for other `Computation` types, such a message is simply discarded.

### Message ID suffixes in Swift {#message-id-suffixes}

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

```cpp
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](https://ytsaurus.tech/docs/en/flow/concepts/swift.md#determinism): the same message ID must still denote the same logical output.

{% note warning %}

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

{% endnote %}

## 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](../concepts/glossary.md#timestamps-and-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](../concepts/glossary.md#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`):

```cpp
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](https://ytsaurus.tech/docs/en/flow/cpp/state.md) section.

## Parameters {#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.

```cpp
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:

```yson
"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 {#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](#parameters)):

```cpp
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](#how-it-works)), and the adjacent `processing_function` field names the function:

```yson
"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:

```cpp
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:

```cpp
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](https://ytsaurus.tech/docs/en/flow/cpp/testing.md) section.

## Example

Full example — [examples/cpp/word_count](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/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](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/examples/cpp/wait_click_join).

## See also

- [Computation (C++)](https://ytsaurus.tech/docs/en/flow/cpp/computation.md)
- [Working with states (C++)](https://ytsaurus.tech/docs/en/flow/cpp/state.md)
- [Testing (C++)](https://ytsaurus.tech/docs/en/flow/cpp/testing.md)
- [Quick start (C++)](https://ytsaurus.tech/docs/en/flow/cpp/getting-started.md)
<!-- endsource: en/_includes/flow/cpp/process-functions.md -->
