Testing in YTsaurus Flow (Python)

Note

This page describes unit testing of Python-pipeline components, as well as integration testing of the full pipeline via FlowTestPythonBase.

General testing architecture

In production, the C++ worker sends gRPC requests to the companion, passing messages, timers, states, and watermarks. 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

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

Test source code

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:

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

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

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

Test source code

YSON State

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

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

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

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

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

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:

Setup

The test inherits from FlowTestPythonBase and sets the PYTHON_COMPANION_BINARY attribute:

Attribute Description
PYTHON_COMPANION_BINARY Path to the Python companion binary

Example of an E2E WordCount test

Warning

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

Use these principles when writing integration tests for pipelines:

  • Test end-to-end scenarios. That is:
  • 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 (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

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.

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:

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

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.

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:

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

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:

ya make -A --test-param 'EXTERNAL_YT_CONFIG={path="//tmp/yt_flow";tablet_cell_bundle=default;proxy_role=default;clusters=[{cluster_name=<cluster-name>};];}'

Speeding up iteration: --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:

ya make -A --ext-py

See also