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
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
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:
- Write input/source data to the local YT.
- In the pipeline spec, mark the sources as
finite=%true. - For a Key Visitor stream 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>"); it switches that stream's dynamicfiniteparameter 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 (a
README.mdis available). - Implement a failover test.
- Many errors surface when workers fail and other workers take over their tasks. So, run a test with
problems=Trueand more than one worker.
- Many errors surface when workers fail and other workers take over their tasks. So, run a test with
- 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.
- Remember that in CI, any part of a test might run unexpectedly long.
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 (.erris the process’s stderr;.logcontains regular logs written viaYT_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 |
|
|
|
Sets the logging level for the process that runs the pipeline (runner). |
|
|
|
|
|
Pauses the test before stopping Flow processes. Combined with |
|
|
(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
- Computation (Python)
- Working with states (Python)
- Distribute flag (Python)
- Testing (Java)
- If you're extending Flow itself — Pipeline testing framework.