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

<!-- source: en/_includes/flow/python/testing.md -->
# Testing in YTsaurus Flow (Python)

{% note info %}

This page describes **unit testing** of Python-[pipeline](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#pipeline) components, as well as **integration testing** of the full pipeline via `FlowTestPythonBase`.

{% endnote %}

## General testing architecture {#architecture}

In production, the C++ [worker](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#worker) sends gRPC requests to the [companion](https://ytsaurus.tech/docs/en/flow/concepts/companion.md), passing [messages](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#message), timers, [states](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#state), and [watermarks](https://ytsaurus.tech/docs/en/flow/concepts/watermarks.md). The companion calls `Computation.do_process()`, which delegates processing to `ProcessFunction`.

In unit tests, pipeline components are tested by directly calling `do_process()` on the `Computation` object with a prepared `RequestContext`.

All tests use the standard pytest.

## Dependencies {#dependencies}

To unit test the companion library, you need to add dependencies to `ya.make`:

```
PEERDIR(
    yt/yt/flow/library/python/companion
)
```

If you’re using Protobuf states, also add a dependency on the proto library.

## Testing ProcessFunction {#testing-process}

[Test source code](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/python/companion/test/test_computation.py)

To test `RowFunction` / `BatchFunction`, you create a `Computation` object and call `do_process()` with a prepared `RequestContext`.

### Creating RequestContext

`RequestContext` is a dataclass that contains the input data for processing:

```python
from yt.yt.flow.library.python.companion.context import RequestContext
from yt.yt.flow.library.python.companion.job import Job
from yt.yt.flow.library.python.companion.stream import (
    StreamIdsMapping,
    StreamSpecs,
    RawStream,
)
from yt.yt.flow.library.python.companion.row import ExtendedMessage


def make_request_ctx(messages=None, timers=None):
    """Create a minimal RequestContext for tests."""
    mapping = StreamIdsMapping({"input": 0, "output": 1})
    specs = StreamSpecs(mapping, [RawStream("input"), RawStream("output")])

    job = Job(
        job_id="test-job",
        computation_id="test-comp",
        stream_specs=specs,
        static_spec={},
    )

    return RequestContext(
        job_id="test-job",
        request_id="test-req",
        computation_id="test-comp",
        messages=messages or [],
        timers=timers or [],
        stream_specs=specs,
        job=job,
    )
```

### Testing RowFunction

```python
from yt.yt.flow.library.python.companion.computation import (
    Computation,
    RowFunction,
)
from yt.yt.flow.library.python.companion.row import (
    ExtendedMessage,
    Message,
)


class PassthroughFunction(RowFunction):
    def on_message(self, message, output, ctx):
        output.add_message(
            Message(message_id=message.message_id, stream_id="output")
        )


def test_passthrough():
    comp = Computation(
        computation_id="test",
        process_function=PassthroughFunction(),
    )
    messages = [
        ExtendedMessage(message_id="m1", stream_id="input"),
        ExtendedMessage(message_id="m2", stream_id="input"),
    ]
    ctx = make_request_ctx(messages=messages)
    response = comp.do_process(ctx)

    assert len(response.transform_results) == 2
    assert response.transform_results[0].messages[0].message_id == "m1"
    assert response.transform_results[1].messages[0].message_id == "m2"
```

### Testing timers

```python
from yt.yt.flow.library.python.companion.row import Timer


class TimerFunction(RowFunction):
    def on_message(self, message, output, ctx):
        output.add_timer(trigger_timestamp=1000, event_timestamp=500)

    def on_timer(self, timer, output, ctx):
        output.add_message(
            Message(message_id="from-timer", stream_id="output")
        )


def test_timer_roundtrip():
    comp = Computation(
        computation_id="test",
        process_function=TimerFunction(),
    )
    messages = [ExtendedMessage(message_id="m1")]
    timers = [Timer(message_id="t1")]
    ctx = make_request_ctx(messages=messages, timers=timers)
    response = comp.do_process(ctx)

    results = response.transform_results
    # The message creates a timer, the timer creates a message
    assert any(r.timers for r in results)
    assert any(r.messages for r in results)
```

## Testing states {#testing-states}

[Test source code](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/python/companion/test/test_context.py)

### YSON State

```python
import yt.type_info as ti

from yt.yt.flow.library.python.companion.context import DefaultRuntimeContext
from yt.yt.flow.library.python.companion.row import (
    ColumnSchema,
    ExtendedMessage,
    Payload,
    TableSchema,
)
from yt.yt.flow.library.python.companion.stream import (
    StreamIdsMapping,
    StreamSpecs,
    RawStream,
)
from yt.yt.flow.library.python.companion.wire_protocol import (
    ColumnValueType,
    UnversionedRow,
    UnversionedValue,
)


KEY_SCHEMA = TableSchema([ColumnSchema("id", ti.String)])


def make_key_payload():
    row = UnversionedRow(values=[
        UnversionedValue(column_id=0, type=ColumnValueType.STRING, value=b"test-key"),
    ])
    return Payload(row, KEY_SCHEMA)


def make_ctx(internal_state_names=None, **kwargs):
    streams = [RawStream("input"), RawStream("output")]
    mapping = StreamIdsMapping({s.stream_id: i for i, s in enumerate(streams)})
    specs = StreamSpecs(mapping, streams)
    return DefaultRuntimeContext(
        internal_state_names=internal_state_names or set(),
        stream_specs=specs,
        internal_states=kwargs.get("internal_states", {}),
        external_states=kwargs.get("external_states", {}),
        watermarks={},
        min_watermark=0,
        computation_parameters={},
        key_schema=KEY_SCHEMA,
    )


def test_yson_state_roundtrip():
    ctx = make_ctx(internal_state_names={"word-state"})
    message = ExtendedMessage(message_id="m1", key=make_key_payload())

    accessor = ctx.state("word-state", message)
    accessor.set({"word": "hello", "count": 1})

    accessor2 = ctx.state("word-state", message)
    result = accessor2.get()
    assert result["word"] == "hello"
    assert result["count"] == 1
```

### Proto State

```python
from google.protobuf.wrappers_pb2 import Int64Value


def test_proto_state_roundtrip():
    ctx = make_ctx(internal_state_names={"join-state"})
    message = ExtendedMessage(message_id="m1", key=make_key_payload())

    accessor = ctx.proto_state("join-state", message, Int64Value)
    state = Int64Value(value=42)
    accessor.set(state)

    accessor2 = ctx.proto_state("join-state", message, Int64Value)
    assert accessor2.get().value == 42
```

### External State

```python
from yt.yt.flow.library.python.companion.state import StatesHolder
from yt.yt.flow.library.python.companion.row import PayloadBuilder


STATE_SCHEMA = TableSchema([
    ColumnSchema("count", ti.Int64),
    ColumnSchema("name", ti.String),
])


def test_external_state_roundtrip():
    ext_holder = StatesHolder("ext", KEY_SCHEMA, STATE_SCHEMA)
    ctx = make_ctx(external_states={"/shuffle-state": ext_holder})
    message = ExtendedMessage(message_id="m1", key=make_key_payload())

    state = ctx.external_state("/shuffle-state", message)
    builder = state.to_builder()
    builder.set("count", 99)
    state.set(builder.finish())

    state2 = ctx.external_state("/shuffle-state", message)
    assert state2.get("count") == 99
```

## Analyzing the response {#analyzing-response}

The `ResponseContext` object returned by `do_process()` contains:

| Field | Type | Description |
|------|------|-------------|
| `transform_results` | `List[TransformResult]` | List of processing results |
| `internal_states` | `Dict[str, StatesHolder]` | Internal states after processing |
| `external_states` | `Dict[str, StatesHolder]` | External states after processing |

Each `TransformResult` contains:

| Field | Type | Description |
|------|------|-------------|
| `parent_ids` | `List[str]` | IDs of parent messages |
| `messages` | `List[Message]` | Output messages |
| `timers` | `List[NewTimer]` | Created timers |

## End-to-end testing with FlowTestPythonBase {#e2e-tests}

Use the `FlowTestPythonBase` base class for full end-to-end pipeline testing (with real C++ workers, queues, and streams).

### Dependencies {#integration-dependencies}

In addition to `PEERDIR` for `integration_test_base`, the integration test needs a cluster recipe, `DEPENDS` for the pipeline binary and `flow_server`, and `DATA` with the spec. The complete `ya.make` for the test from [WordCount](https://ytsaurus.tech/docs/en/flow/python/examples/wordcount.md):

{% code '/yt/yt/flow/examples/python/word_count/test/ya.make' lang='text' %}

### Setup {#python-test-setup}

The test inherits from `FlowTestPythonBase` and sets the `PYTHON_COMPANION_BINARY` attribute:

{% code '/yt/yt/flow/examples/python/word_count/test/test_wordcount.py' lang='python' lines='[BEGIN test_setup]-[END test_setup]' %}

| Attribute | Description |
|----------|-------------|
| `PYTHON_COMPANION_BINARY` | Path to the Python companion binary |

[Example of an E2E WordCount test](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/examples/python/word_count/test/test_wordcount.py)


{% note warning %}

Integration tests require a deployed YTsaurus cluster and are run with `ya make -tt`. For fast iteration, use the unit tests described above.

{% endnote %}

<!-- source: en/_includes/flow/testing-integration-body.md -->
Use these principles when writing integration tests for pipelines:

* Test end-to-end scenarios. That is:
  * Write input/source data to the local YT.
  * In the pipeline spec, mark the sources as `finite=%true`.
  * For a [Key Visitor stream](https://ytsaurus.tech/docs/en/flow/concepts/key_visitor.md) that should run during the test, set `finite=%false`.
  * Run the pipeline.
  * When the Key Visitor should finish, call the integration test base method [`self.ask_key_visitor_to_complete("<computation_id>", "<stream_id>")`](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/python/integration_test_base/yt_flow_base.py); it switches that stream's [dynamic `finite` parameter](https://ytsaurus.tech/docs/en/flow/concepts/key_visitor.md#dynamic-params) to `%true`.
  * Wait for the pipeline to finish.
  * Check the output data.
* Prepare the environment with the same code you use in production:
  * Test the same pipeline binary that will run in production.
  * Generate pipeline specs with the same code that generates them for production.
  * If input/output data require non-trivial serialization/parsing, implement this logic with code shared with production.
* Use the [shared test framework](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/python/integration_test_base) (a `README.md` is available).
* Implement a failover test.
  * Many errors surface when workers fail and other workers take over their tasks. So, run a test with `problems=True` and more than one worker.
* Write stable tests.
  * Remember that in CI, any part of a test might run unexpectedly long.
    * The ideal runtime for a single test is 20 seconds for any build type.
    * For sanitizer builds, reduce the amount of input data.
    * Set all local timeouts with a multiple safety margin.
  * The test execution logic shouldn’t depend significantly on the current time. For example, a test shouldn’t fail if it starts before midnight and finishes after.


## Debugging tests {#debug}

### Logs {#debug-logs}

After a test finishes, you can find its logs in `test-results/py3test/testing_out_stuff` in the test directory. Key logs include:

* `run.log` — Python test logs.
* `<test_class_name>/<test_name>/Controller_<number>...` — Controller logs (`.err` is the process’s stderr; `.log` contains regular logs written via `YT_LOG_...`).
* `<test_class_name>/<test_name>/Worker_<number>...` — Worker logs.
* `<test_class_name>/<test_name>/Runner...` — Runner logs.

If something isn’t working, check errors in all these logs. The approach to reviewing them is the same as in production — see [Logs](https://ytsaurus.tech/docs/en/flow/devops/vanilla/diagnostics/logs.md).

You can also view logs before the test finishes. To do this, locate the temporary directory where the test runs. The simplest way is to run the test with the `--keep-temps` flag: `ya make --keep-temps -ttt <target>`. In this case, `ya make` won’t delete the temporary directory after the test finishes and will print a link to it in the output.

You can also use these bash aliases:

```bash
alias curtestdir="ps -f -u $USER | python3 -c \"import sys, re; drs = set(e for e in re.findall(r'[\s=](/\S*testing' + r'_out_stuff)\b', sys.stdin.read())); print('' if len(drs) == 1 else 'Select first from ' + repr(drs), file=sys.stderr); print(list(drs)[0])\""
alias cdcurtestdir='cd $(curtestdir)'

# Get the link to the local YTsaurus UI from the logs.
alias curlocalyt='cat $(curtestdir)/stderr 2>/dev/null | grep YT'
```


### Test framework {#debug-test-framework}

You can read about how to set up the environment for better test framework performance and how to influence test parameters in the framework’s [README.md](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/python/integration_test_base).
<!-- endsource: en/_includes/flow/testing-integration-body.md -->

<!-- source: en/_includes/flow/testing-test-param-body.md -->
You configure the behavior of integration tests using `--test-param NAME=VAL`:

#|
|| **Parameter** | **Default** | **Values** | **Action** ||
|| `RUNNER_LOG_LEVEL` | | `Error`, `Info`, `Debug`, … | Sets the logging level for the process that runs the pipeline (runner). ||
|| `PAUSE_BEFORE_FLOW_PROCESS_FEDERATION_TEARDOWN` | `0` | `0`, `1` | Pauses the test before stopping Flow processes. Combined with `--test-disable-timeout`, this lets you keep the local YTsaurus and the Flow process federation running for a long time so you can study them at your own pace via the UI. ||
|| `EXTERNAL_YT_CONFIG` | (not set) | YSON — see below | Runs the pipeline on real external YTsaurus clusters instead of the local recipe. ||
|#

Examples:

```bash
ya make -A --test-param RUNNER_LOG_LEVEL=Debug
ya make -A --test-disable-timeout --test-param PAUSE_BEFORE_FLOW_PROCESS_FEDERATION_TEARDOWN=1
```

### `EXTERNAL_YT_CONFIG` {#external-yt-config}

The local YTsaurus still starts via the recipe, but the test ignores it.

Required fields that are common for all clusters:

- `path` — base directory.
- `tablet_cell_bundle` — bundle for the dynamic tables that are created.
- `proxy_role` — RPC proxy role.

Optional: `primary_medium` (default `"default"`).

The `clusters` list — the first element is primary. Record fields: `cluster_name` (required), `proxy_url` (defaults to `cluster_name`).

Authorization: when using an external YTsaurus, `YT_TOKEN`/`YT_USER` from the local recipe are cleared; yt-wrapper picks up the token from `~/.yt/token`.

Isolation: `work_yt_path = path/<local-username>/<test_name>`; the `path/<username>` directory is deleted and recreated once per class in `setup_class`.

Example:

```bash
ya make -A --test-param 'EXTERNAL_YT_CONFIG={path="//tmp/yt_flow";tablet_cell_bundle=default;proxy_role=default;clusters=[{cluster_name=<cluster-name>};];}'
```
<!-- endsource: en/_includes/flow/testing-test-param-body.md -->

## Speeding up iteration: `--ext-py` {#ext-py}

The `--ext-py` flag avoids linking the binary on every `ya make` run if the changes from the previous run affected only `*.py` files:

```bash
ya make -A --ext-py
```

## See also

- [Computation (Python)](https://ytsaurus.tech/docs/en/flow/python/computation.md)
- [Working with states (Python)](https://ytsaurus.tech/docs/en/flow/python/state.md)
- [Distribute flag (Python)](https://ytsaurus.tech/docs/en/flow/python/distribute.md)
- [Testing (Java)](https://ytsaurus.tech/docs/en/flow/java/testing.md)
- If you're extending Flow itself — [Pipeline testing framework](https://ytsaurus.tech/docs/en/flow/contributor/testing-framework.md).
<!-- endsource: en/_includes/flow/python/testing.md -->
