Тестирование с TestComputationHarness в YTsaurus Flow (Java)

Модуль flow-test-utils предоставляет утилиты для юнит-тестирования компонентов Flow-пайплайна на Java и Kotlin без запуска реального gRPC-сервера и C++ воркеров. Центральный класс — TestComputationHarness, который эмулирует вызовы doProcess на уровне компаньона.

Примечание

Класс TestComputationHarness предназначен для юнит-тестирования отдельных ProcessFunction без запуска реального пайплайна. Для интеграционного тестирования полного пайплайна (с реальными C++ воркерами, очередями и стримами) используйте FlowTestJavaBase — см. раздел ниже.

Общая архитектура тестирования

В продакшене C++ воркер отправляет gRPC-запросы компаньону, передавая сообщения, таймеры, состояние и watermark-и. Companion вызывает Computation.doProcess(), который делегирует обработку в ProcessFunction.

В тестах TestComputationHarness заменяет gRPC-слой: он принимает TestDoProcessRequest, конвертирует его в protobuf-формат, вызывает CompanionRequestProcessor.processBatch() и возвращает TestDoProcessResponse с десериализованными результатами.

Зависимости

Для использования тестовых утилит необходимо добавить зависимость на flow-test-utils:

PEERDIR(
    yt/java/flow/flow-test-utils
)

Все примеры ниже написаны с использованием JUnit версии 5+.

Настройка TestComputationHarness

Для создания TestComputationHarness необходимо:

  1. PipelineContext — контекст пайплайна с зарегистрированными объектами Computation и стримами.
  2. Pipeline spec — статическая спецификация пайплайна в формате YSON (файл pipeline.yson).
  3. External state schemas (опционально) — схемы внешних состояний, если computation использует ExternalStateAccessor.

Builder API

TestComputationHarness harness = TestComputationHarness.builder()
        .setPipelineContext(pipelineContext)       // обязательно
        .setPipelineSpec(txtSpec)                  // обязательно: String, YTreeNode или InputStream
        .addExternalStateSchema("state-name", schema) // опционально
        .setJobContext(jobContext)                 // опционально, по умолчанию timeout=10 мин
        .build();
val harness = TestComputationHarness.builder()
        .setPipelineContext(pipelineContext)       // обязательно
        .setPipelineSpec(txtSpec)                  // обязательно: String, YTreeNode или InputStream
        .addExternalStateSchema("state-name", schema) // опционально
        .setJobContext(jobContext)                 // опционально, по умолчанию timeout=10 мин
        .build()

При вызове build() harness автоматически:

  • Извлекает информацию о стримах из pipeline spec.
  • Регистрирует недостающие стримы как нетипизированные (через FlowStreams.raw) в PipelineContext.
  • Создаёт внутренний CompanionRequestProcessor.

Построение тестовых запросов

TestDoProcessRequest

TestDoProcessRequest — запрос на обработку батча сообщений и/или таймеров. Создаётся через builder:

var request = TestDoProcessRequest.builder("computation-id")
        .setMessages(messages)           // List<ExtendedMessage>
        .setTimers(timers)               // List<Timer>
        .setWatermarks(watermarks)       // Map<String, Long>
        .setExternalState("name", stateMap)  // Map<Payload, ExternalState>
        .setInternalState("name", stateMap)  // Map<Payload, InternalState>
        .build();
val request = TestDoProcessRequest.builder("computation-id")
        .setMessages(messages)           // List<ExtendedMessage>
        .setTimers(timers)               // List<Timer>
        .setWatermarks(watermarks)       // Map<String, Long>
        .setExternalState("name", stateMap)  // Map<Payload, ExternalState>
        .setInternalState("name", stateMap)  // Map<Payload, InternalState>
        .build()

Параметры:

  • computationId — идентификатор computation, которому адресован запрос.
  • messages — список входных сообщений типа ExtendedMessage.
  • timers — список сработавших таймеров типа Timer.
  • watermarks — карта watermark-ов по стримам (streamId → timestamp).
  • externalStates — предзаполненное внешнее состояние.
  • internalStates — предзаполненное внутреннее состояние.

Построение сообщений

С типизированными стримами

Если стримы зарегистрированы как типизированные (через FlowStreams.typed(...)), payload можно передавать как POJO-объект:

ExtendedMessage msg = ExtendedMessage.builder()
        .setStreamId("hit")
        .setKey(key)
        .setPayload(new Hit("hit-1", 1000L, "payload-data"))
        .setEventTimestamp(1000L)
        .build();
val msg = ExtendedMessage.builder()
        .setStreamId("hit")
        .setKey(key)
        .setPayload(Hit("hit-1", 1000L, "payload-data"))
        .setEventTimestamp(1000L)
        .build()

С нетипизированными стримами (PayloadBuilder)

Если стримы не типизированы, payload строится через PayloadBuilder:

var hitSchema = testHarness.getStream("hit").getSchema();
var payload = new PayloadBuilder(hitSchema)
        .set("hit_id", "hit-1")
        .set("hit_time", 1000L)
        .set("hit_payload", "payload-data")
        .finish();

ExtendedMessage msg = ExtendedMessage.builder()
        .setStreamId("hit")
        .setKey(key)
        .setPayload(payload)
        .setEventTimestamp(1000L)
        .build();
val hitSchema = testHarness.getStream("hit").getSchema()
val payload = PayloadBuilder(hitSchema)
        .set("hit_id", "hit-1")
        .set("hit_time", 1000L)
        .set("hit_payload", "payload-data")
        .finish()

val msg = ExtendedMessage.builder()
        .setStreamId("hit")
        .setKey(key)
        .setPayload(payload)
        .setEventTimestamp(1000L)
        .build()

Примечание

Метод testHarness.getStream(streamId) возвращает FlowStream<?> со схемой, извлечённой из pipeline spec. Это удобно для получения схемы стрима без ручного определения.

Построение ключей

Ключ (Payload) строится по схеме group_by_schema тестируемого Computation из pipeline spec.

Примечание

Метод testHarness.getGroupBySchema(computationId) возвращает TableSchema на основе group_by_schema, извлечённой из pipeline spec для тестируемого Computation.

TableSchema keySchema = testHarness.getGroupBySchema("join");

Payload key = new PayloadBuilder(keySchema)
        .set("hash", 0L)
        .set("hit_id", "hit-1")
        .set("hit_time", 1000L)
        .finish();
val keySchema = testHarness.getGroupBySchema("join")

val key = PayloadBuilder(keySchema)
        .set("hash", 0L)
        .set("hit_id", "hit-1")
        .set("hit_time", 1000L)
        .finish()

Совет

Поле hash (farm_hash) вычисляется на стороне C++ воркера. В тестах его можно устанавливать в 0L.

Построение таймеров

Таймер создаётся через builder:

Timer timer = Timer.builder()
        .setStreamId("timer")
        .setKey(key)
        .setTriggerTimestamp(triggerTimestamp)
        .build();
val timer = Timer.builder()
        .setStreamId("timer")
        .setKey(key)
        .setTriggerTimestamp(triggerTimestamp)
        .build()

Настройка watermark-ов

Watermark-и определяют, какие сообщения считаются «опоздавшими» (late). Сообщение с eventTimestamp < watermark будет отброшено, если в ProcessFunction реализована соответствующая проверка через ctx.getEpochInputEventWatermark(). Пример такой проверки в Wait Click Join на Java.

Map<String, Long> watermarks = Map.of(
        "hit", 0L,
        "action", 0L
);
val watermarks = mapOf(
        "hit" to 0L,
        "action" to 0L
)

Установка watermark в 0L означает, что все сообщения с eventTimestamp >= 0 не будут считаться опоздавшими.

Предзаполнение состояния

External State

TableSchema stateSchema = TableSchema.builder()
        .addValue("hit_payload", ColumnValueType.STRING)
        .addValue("show_time", ColumnValueType.UINT64)
        .addValue("click_time", ColumnValueType.UINT64)
        .build();

ExternalState preState = new ExternalState(
        new PayloadBuilder(stateSchema)
                .set("hit_payload", "some-payload")
                .set("show_time", showTime)
                .set("click_time", clickTime)
                .finish()
);

var request = TestDoProcessRequest.builder("join")
        .setTimers(List.of(timer))
        .setExternalState("join-state", Map.of(key, preState))
        .build();
val stateSchema = TableSchema.builder()
        .addValue("hit_payload", ColumnValueType.STRING)
        .addValue("show_time", ColumnValueType.UINT64)
        .addValue("click_time", ColumnValueType.UINT64)
        .build()

val preState = ExternalState(
        PayloadBuilder(stateSchema)
                .set("hit_payload", "some-payload")
                .set("show_time", showTime)
                .set("click_time", clickTime)
                .finish()
)

val request = TestDoProcessRequest.builder("join")
        .setTimers(listOf(timer))
        .setExternalState("join-state", mapOf(key to preState))
        .build()

Internal State (Proto)

Для внутреннего состояния на основе protobuf:

TJoinState protoState = TJoinState.newBuilder()
        .setHitPayload("some-payload")
        .setShowTime(showTime)
        .setClickTime(clickTime)
        .build();

InternalState internalState = new InternalState(protoState.toByteArray());

var request = TestDoProcessRequest.builder("join")
        .setTimers(List.of(timer))
        .setInternalState("join-state", Map.of(key, internalState))
        .build();
val protoState = TJoinState.newBuilder()
        .setHitPayload("some-payload")
        .setShowTime(showTime)
        .setClickTime(clickTime)
        .build()

val internalState = InternalState(protoState.toByteArray())

val request = TestDoProcessRequest.builder("join")
        .setTimers(listOf(timer))
        .setInternalState("join-state", mapOf(key to internalState))
        .build()

Анализ результатов

TestDoProcessResponse

TestDoProcessResponse предоставляет методы для проверки результатов обработки.

Метод Описание
getOutputMessagesFlatten() Все выходные сообщения из всех transform-групп
getOutputTimersFlatten() Все установленные таймеры из всех transform-групп
getTransformResults() Список TransformResult с группировкой по parent ids

Примечание

TestDoProcessResponse предоставляет два набора методов доступа к состоянию:

  • Полная картина (getExternalState(s), getInternalState(s)) — все состояния после обработки: ключи, переданные в запросе, с наложенными поверх изменениями computation. Используйте, когда нужно проверить итоговое состояние независимо от того, изменил ли его computation.
  • Только изменённые (getModifiedExternalState(s), getModifiedInternalState(s)) — только то, что computation действительно изменил. Именно эти состояния отправляются обратно воркеру и сохраняются. Используйте для проверки того, что стейт по ключу не менялся, или когда нужно проверить только изменения состояния.

Все методы доступа к состоянию null-safe: неизвестное имя состояния или ключ возвращают пустую карту или null, а не выбрасывают исключение.

Полная картина состояния (загруженные состояния с наложенными изменениями computation):

Метод Описание
getExternalStateNames() / getInternalStateNames() Имена всех состояний (загруженных и/или изменённых)
getExternalStates() / getInternalStates() Все состояния, сгруппированные по имени
getExternalStates(name) / getInternalStates(name) Все записи состояния по имени (пустая карта, если нет)
getExternalStateSize(name) / getInternalStateSize(name) Количество записей в состоянии (0, если нет)
getExternalState(name, key) / getInternalState(name, key) Конкретная запись состояния (null, если нет)

Только состояния, изменённые computation:

Метод Описание
getModifiedExternalStateNames() / getModifiedInternalStateNames() Имена изменённых состояний
getModifiedExternalStates() / getModifiedInternalStates() Изменённые состояния, сгруппированные по имени
getModifiedExternalStates(name) / getModifiedInternalStates(name) Изменённые записи состояния по имени (пустая карта, если нет)
getModifiedExternalStateSize(name) / getModifiedInternalStateSize(name) Количество изменённых записей (0, если нет)
getModifiedExternalState(name, key) / getModifiedInternalState(name, key) Конкретная изменённая запись (null, если не изменена)

Проверка выходных сообщений

var response = testHarness.doProcess(request);

// Проверяем количество выходных сообщений
assertEquals(1, response.getOutputMessagesFlatten().size());

// Получаем сообщение и проверяем стрим
var msg = response.getOutputMessagesFlatten().get(0);
assertEquals("joined_action", msg.getStreamId());

// Для типизированных стримов — приведение payload к POJO
JoinedAction result = (JoinedAction) msg.getPayload();
assertEquals("hit-1", result.getHitId());

// Для нетипизированных стримов — чтение полей через get()
byte[] data = msg.get("data", byte[].class);
val response = testHarness.doProcess(request)

// Проверяем количество выходных сообщений
assertEquals(1, response.getOutputMessagesFlatten().size)

// Получаем сообщение и проверяем стрим
val msg = response.getOutputMessagesFlatten()[0]
assertEquals("joined_action", msg.getStreamId())

// Для типизированных стримов — приведение payload к POJO
val result = msg.getPayload() as JoinedAction
assertEquals("hit-1", result.getHitId())

// Для нетипизированных стримов — чтение полей через get()
val data = msg.get("data", ByteArray::class.java)

Проверка таймеров

assertEquals(1, response.getOutputTimersFlatten().size());
var timer = response.getOutputTimersFlatten().get(0);
assertEquals(expectedTriggerTimestamp, timer.getTriggerTimestamp());
assertEquals(1, response.getOutputTimersFlatten().size)
val timer = response.getOutputTimersFlatten()[0]
assertEquals(expectedTriggerTimestamp, timer.getTriggerTimestamp())

Проверка состояния

// External state
assertEquals(1, response.getExternalStateSize("join-state"));
var state = response.getExternalState("join-state", key);
assertFalse(state.isReset());
assertEquals("payload-data", state.getValue().get("hit_payload", String.class));

// Проверка сброса состояния
var stateAfterTimer = response.getExternalState("join-state", key);
assertTrue(stateAfterTimer.isReset());
// External state
assertEquals(1, response.getExternalStateSize("join-state"))
val state = response.getExternalState("join-state", key)!!
assertFalse(state.isReset)
assertEquals("payload-data", state.value.get("hit_payload", String::class.java))

// Проверка сброса состояния
val stateAfterTimer = response.getExternalState("join-state", key)!!
assertTrue(stateAfterTimer.isReset)
// Internal state (Proto)
var internalState = response.getInternalState("join-state", key);
assertFalse(internalState.isReset());
TJoinState joinState = ProtoUtils.parseBytes(internalState.getValue(), TJoinState.class);
assertEquals("payload-data", joinState.getHitPayload());
// Internal state (Proto)
val internalState = response.getInternalState("join-state", key)!!
assertFalse(internalState.isReset)
val joinState = ProtoUtils.parseBytes(internalState.value, TJoinState::class.java)
assertEquals("payload-data", joinState.getHitPayload())

Полная картина и изменённые состояния

Методы getExternalState / getInternalState возвращают полную картину: предзаполненное состояние доступно, даже если computation его не менял. Чтобы проверить, что именно изменил computation, используйте методы getModified*.

// Полная картина: загруженное состояние видно, даже если не изменялось
assertEquals(1, response.getExternalStateSize("join-state"));
assertNotNull(response.getExternalState("join-state", key));

// Только изменённые состояния: computation ничего не записал для этого ключа
assertEquals(0, response.getModifiedExternalStateSize("join-state"));
assertNull(response.getModifiedExternalState("join-state", key));

// Неизвестные имя/ключ не бросают исключение
assertNull(response.getExternalState("unknown", key));
assertTrue(response.getExternalStates("unknown").isEmpty());
// Полная картина: загруженное состояние видно, даже если не изменялось
assertEquals(1, response.getExternalStateSize("join-state"))
assertNotNull(response.getExternalState("join-state", key))

// Только изменённые состояния: computation ничего не записал для этого ключа
assertEquals(0, response.getModifiedExternalStateSize("join-state"))
assertNull(response.getModifiedExternalState("join-state", key))

// Неизвестные имя/ключ не бросают исключение
assertNull(response.getExternalState("unknown", key))
assertTrue(response.getExternalStates("unknown").isEmpty())

Пример: тест без Spring

Полный пример юнит-теста для JoinProcessFunction из проекта wait_click_join. В этом варианте PipelineContext создаётся вручную, без Spring-контейнера.

public class JoinProcessFunctionTest {

    private static final long WAIT_SECONDS = 10L;
    private static final long BASE_HIT_TIME = 1000L;
    private static final long WATERMARK = 0L;

    private TestComputationHarness testHarness;
    private TableSchema keySchema;
    private TableSchema joinStateSchema;

    @BeforeEach
    void init() throws IOException {
        // 1. Создаём PipelineContext вручную
        var pipelineContext = new PipelineContext();

        // 2. Регистрируем computation с process function
        Computation join = Computation.builder()
                .setComputationId("join")
                .setProcessFunction(new JoinProcessFunction())
                .build();
        pipelineContext.registerComputation(join);

        // 3. Регистрируем типизированные стримы
        pipelineContext.registerStream(FlowStreams.typed("hit", Hit.class));
        pipelineContext.registerStream(FlowStreams.typed("action", Action.class));
        pipelineContext.registerStream(FlowStreams.typed("joined_action", JoinedAction.class));

        // 4. Определяем схему внешнего состояния
        this.joinStateSchema = TableSchema.builder()
                .addValue("hit_payload", ColumnValueType.STRING)
                .addValue("show_time", ColumnValueType.UINT64)
                .addValue("click_time", ColumnValueType.UINT64)
                .build();

        // 5. Читаем pipeline spec и создаём harness
        var specPath = Paths.getSourcePath(
                "yt/yt/flow/examples/java/wait_click_join/test/pipeline.yson");
        var txtSpec = Files.readString(Path.of(specPath));

        this.testHarness = TestComputationHarness.builder()
                .setPipelineContext(pipelineContext)
                .setPipelineSpec(txtSpec)
                .addExternalStateSchema("join-state", joinStateSchema)
                .build();

        // 6. Получаем key schema для join computation (group_by_schema из pipeline.yson)
        this.keySchema = testHarness.getGroupBySchema("join");
    }

    // — Вспомогательные методы —

    private Payload buildKey(String hitId, long hitTime) {
        return new PayloadBuilder(keySchema)
                .set("hash", 0L)
                .set("hit_id", hitId)
                .set("hit_time", hitTime)
                .finish();
    }

    private ExtendedMessage buildHitMessage(String hitId, long hitTime, String hitPayload) {
        return ExtendedMessage.builder()
                .setStreamId("hit")
                .setKey(buildKey(hitId, hitTime))
                .setPayload(new Hit(hitId, hitTime, hitPayload))
                .setEventTimestamp(hitTime)
                .build();
    }

    private Map<String, Long> defaultWatermarks() {
        return Map.of("hit", WATERMARK, "action", WATERMARK);
    }

    // — Тесты —

    @Test
    void testHitMessageStoresHitPayloadInState() {
        var messages = List.of(buildHitMessage("hit-1", BASE_HIT_TIME, "payload-data"));

        var request = TestDoProcessRequest.builder("join")
                .setMessages(messages)
                .setWatermarks(defaultWatermarks())
                .build();

        var response = testHarness.doProcess(request);

        // Нет выходных сообщений — результат будет при срабатывании таймера
        assertTrue(response.getOutputMessagesFlatten().isEmpty());

        // Таймер должен быть установлен
        assertEquals(1, response.getOutputTimersFlatten().size());
        assertEquals(BASE_HIT_TIME + WAIT_SECONDS,
                response.getOutputTimersFlatten().get(0).getTriggerTimestamp());

        // Состояние должно содержать hit_payload
        Payload key = buildKey("hit-1", BASE_HIT_TIME);
        var state = response.getExternalState("join-state", key);
        assertFalse(state.isReset());
        assertEquals("payload-data", state.getValue().get("hit_payload", String.class));
    }

    @Test
    void testTimerEmitsJoinedAction() {
        Payload key = buildKey("hit-10", BASE_HIT_TIME);

        // Предзаполняем состояние
        ExternalState preState = new ExternalState(
                new PayloadBuilder(joinStateSchema)
                        .set("hit_payload", "some-payload")
                        .set("show_time", BASE_HIT_TIME + 3L)
                        .set("click_time", BASE_HIT_TIME + 7L)
                        .finish()
        );

        Timer timer = Timer.builder()
                .setStreamId("timer")
                .setKey(key)
                .setTriggerTimestamp(BASE_HIT_TIME + WAIT_SECONDS)
                .build();

        var request = TestDoProcessRequest.builder("join")
                .setTimers(List.of(timer))
                .setExternalState("join-state", Map.of(key, preState))
                .build();

        var response = testHarness.doProcess(request);

        // Должно быть одно выходное сообщение
        assertEquals(1, response.getOutputMessagesFlatten().size());
        var msg = response.getOutputMessagesFlatten().get(0);
        assertEquals("joined_action", msg.getStreamId());

        JoinedAction result = (JoinedAction) msg.getPayload();
        assertEquals("hit-10", result.getHitId());
        assertTrue(result.getClick());

        // Состояние должно быть сброшено
        assertTrue(response.getExternalState("join-state", key).isReset());
    }

    @Test
    void testLateMessageIsDropped() {
        long watermark = BASE_HIT_TIME + 5L;
        long eventTimestamp = BASE_HIT_TIME + 3L; // < watermark → late

        var messages = List.of(buildActionMessage("hit-5", BASE_HIT_TIME,
                BASE_HIT_TIME + 2L, false, eventTimestamp));

        var request = TestDoProcessRequest.builder("join")
                .setMessages(messages)
                .setWatermarks(Map.of("hit", watermark, "action", watermark))
                .build();

        var response = testHarness.doProcess(request);

        // Опоздавшее сообщение отброшено
        assertTrue(response.getOutputMessagesFlatten().isEmpty());
        assertTrue(response.getOutputTimersFlatten().isEmpty());
        assertEquals(0, response.getExternalStateSize("join-state"));
    }
}
class JoinProcessFunctionTest {

    companion object {
        private const val WAIT_SECONDS = 10L
        private const val BASE_HIT_TIME = 1000L
        private const val WATERMARK = 0L
    }

    private lateinit var testHarness: TestComputationHarness
    private lateinit var keySchema: TableSchema
    private lateinit var joinStateSchema: TableSchema

    @BeforeEach
    fun init() {
        // 1. Создаём PipelineContext вручную
        val pipelineContext = PipelineContext()

        // 2. Регистрируем computation с process function
        val join = Computation.builder()
                .setComputationId("join")
                .setProcessFunction(JoinProcessFunction())
                .build()
        pipelineContext.registerComputation(join)

        // 3. Регистрируем типизированные стримы
        pipelineContext.registerStream(FlowStreams.typed("hit", Hit::class.java))
        pipelineContext.registerStream(FlowStreams.typed("action", Action::class.java))
        pipelineContext.registerStream(FlowStreams.typed("joined_action", JoinedAction::class.java))

        // 4. Определяем схему внешнего состояния
        joinStateSchema = TableSchema.builder()
                .addValue("hit_payload", ColumnValueType.STRING)
                .addValue("show_time", ColumnValueType.UINT64)
                .addValue("click_time", ColumnValueType.UINT64)
                .build()

        // 5. Читаем pipeline spec и создаём harness
        val specPath = Paths.getSourcePath(
                "yt/yt/flow/examples/java/wait_click_join/test/pipeline.yson")
        val txtSpec = Files.readString(Path.of(specPath))

        testHarness = TestComputationHarness.builder()
                .setPipelineContext(pipelineContext)
                .setPipelineSpec(txtSpec)
                .addExternalStateSchema("join-state", joinStateSchema)
                .build()

        // 6. Получаем key schema для join computation (group_by_schema из pipeline.yson)
        keySchema = testHarness.getGroupBySchema("join")
    }

    // — Вспомогательные методы —

    private fun buildKey(hitId: String, hitTime: Long): Payload =
        PayloadBuilder(keySchema)
                .set("hash", 0L)
                .set("hit_id", hitId)
                .set("hit_time", hitTime)
                .finish()

    private fun buildHitMessage(hitId: String, hitTime: Long, hitPayload: String): ExtendedMessage =
        ExtendedMessage.builder()
                .setStreamId("hit")
                .setKey(buildKey(hitId, hitTime))
                .setPayload(Hit(hitId, hitTime, hitPayload))
                .setEventTimestamp(hitTime)
                .build()

    private fun defaultWatermarks() = mapOf("hit" to WATERMARK, "action" to WATERMARK)

    // — Тесты —

    @Test
    fun testHitMessageStoresHitPayloadInState() {
        val messages = listOf(buildHitMessage("hit-1", BASE_HIT_TIME, "payload-data"))

        val request = TestDoProcessRequest.builder("join")
                .setMessages(messages)
                .setWatermarks(defaultWatermarks())
                .build()

        val response = testHarness.doProcess(request)

        // Нет выходных сообщений — результат будет при срабатывании таймера
        assertTrue(response.getOutputMessagesFlatten().isEmpty())

        // Таймер должен быть установлен
        assertEquals(1, response.getOutputTimersFlatten().size)
        assertEquals(BASE_HIT_TIME + WAIT_SECONDS,
                response.getOutputTimersFlatten()[0].getTriggerTimestamp())

        // Состояние должно содержать hit_payload
        val key = buildKey("hit-1", BASE_HIT_TIME)
        val state = response.getExternalState("join-state", key)!!
        assertFalse(state.isReset)
        assertEquals("payload-data", state.value.get("hit_payload", String::class.java))
    }

    @Test
    fun testTimerEmitsJoinedAction() {
        val key = buildKey("hit-10", BASE_HIT_TIME)

        // Предзаполняем состояние
        val preState = ExternalState(
                PayloadBuilder(joinStateSchema)
                        .set("hit_payload", "some-payload")
                        .set("show_time", BASE_HIT_TIME + 3L)
                        .set("click_time", BASE_HIT_TIME + 7L)
                        .finish()
        )

        val timer = Timer.builder()
                .setStreamId("timer")
                .setKey(key)
                .setTriggerTimestamp(BASE_HIT_TIME + WAIT_SECONDS)
                .build()

        val request = TestDoProcessRequest.builder("join")
                .setTimers(listOf(timer))
                .setExternalState("join-state", mapOf(key to preState))
                .build()

        val response = testHarness.doProcess(request)

        // Должно быть одно выходное сообщение
        assertEquals(1, response.getOutputMessagesFlatten().size)
        val msg = response.getOutputMessagesFlatten()[0]
        assertEquals("joined_action", msg.getStreamId())

        val result = msg.getPayload() as JoinedAction
        assertEquals("hit-10", result.getHitId())
        assertTrue(result.getClick())

        // Состояние должно быть сброшено
        assertTrue(response.getExternalState("join-state", key)!!.isReset)
    }

    @Test
    fun testLateMessageIsDropped() {
        val watermark = BASE_HIT_TIME + 5L
        val eventTimestamp = BASE_HIT_TIME + 3L // < watermark → late

        val messages = listOf(buildActionMessage("hit-5", BASE_HIT_TIME,
                BASE_HIT_TIME + 2L, false, eventTimestamp))

        val request = TestDoProcessRequest.builder("join")
                .setMessages(messages)
                .setWatermarks(mapOf("hit" to watermark, "action" to watermark))
                .build()

        val response = testHarness.doProcess(request)

        // Опоздавшее сообщение отброшено
        assertTrue(response.getOutputMessagesFlatten().isEmpty())
        assertTrue(response.getOutputTimersFlatten().isEmpty())
        assertEquals(0, response.getExternalStateSize("join-state"))
    }
}

Многошаговый тест

Для тестирования полного потока обработки (hit → show → click → timer) состояние передаётся между шагами вручную:

@Test
void testFullJoinFlow() {
    String hitId = "hit-20";
    long hitTime = BASE_HIT_TIME;

    // Шаг 1: обработка hit-сообщения
    var hitRequest = TestDoProcessRequest.builder("join")
            .setMessages(List.of(buildHitMessage(hitId, hitTime, "full-payload")))
            .setWatermarks(defaultWatermarks())
            .build();
    var hitResponse = testHarness.doProcess(hitRequest);

    Payload key = buildKey(hitId, hitTime);
    var stateAfterHit = hitResponse.getExternalState("join-state", key);

    // Шаг 2: обработка show-сообщения, передаём состояние из шага 1
    var showRequest = TestDoProcessRequest.builder("join")
            .setMessages(List.of(buildActionMessage(hitId, hitTime, hitTime + 3L, false, hitTime + 1L)))
            .setExternalState("join-state", Map.of(key, stateAfterHit))
            .setWatermarks(defaultWatermarks())
            .build();
    var showResponse = testHarness.doProcess(showRequest);

    var stateAfterShow = showResponse.getExternalState("join-state", key);

    // Шаг 3: обработка click-сообщения
    var clickRequest = TestDoProcessRequest.builder("join")
            .setMessages(List.of(buildActionMessage(hitId, hitTime, hitTime + 7L, true, hitTime + 2L)))
            .setExternalState("join-state", Map.of(key, stateAfterShow))
            .setWatermarks(defaultWatermarks())
            .build();
    var clickResponse = testHarness.doProcess(clickRequest);

    var stateAfterClick = clickResponse.getExternalState("join-state", key);

    // Шаг 4: срабатывание таймера
    Timer timer = Timer.builder()
            .setStreamId("timer")
            .setKey(key)
            .setTriggerTimestamp(hitTime + WAIT_SECONDS)
            .build();
    var timerRequest = TestDoProcessRequest.builder("join")
            .setTimers(List.of(timer))
            .setExternalState("join-state", Map.of(key, stateAfterClick))
            .build();
    var timerResponse = testHarness.doProcess(timerRequest);

    // Проверяем финальный результат
    assertEquals(1, timerResponse.getOutputMessagesFlatten().size());
    JoinedAction result = timerResponse.getOutputMessagesFlatten().get(0).getPayload();
    assertEquals(hitId, result.getHitId());
    assertTrue(result.getClick());

    assertTrue(timerResponse.getExternalState("join-state", key).isReset());
}
@Test
fun testFullJoinFlow() {
    val hitId = "hit-20"
    val hitTime = BASE_HIT_TIME

    // Шаг 1: обработка hit-сообщения
    val hitRequest = TestDoProcessRequest.builder("join")
            .setMessages(listOf(buildHitMessage(hitId, hitTime, "full-payload")))
            .setWatermarks(defaultWatermarks())
            .build()
    val hitResponse = testHarness.doProcess(hitRequest)

    val key = buildKey(hitId, hitTime)
    val stateAfterHit = hitResponse.getExternalState("join-state", key)!!

    // Шаг 2: обработка show-сообщения, передаём состояние из шага 1
    val showRequest = TestDoProcessRequest.builder("join")
            .setMessages(listOf(buildActionMessage(hitId, hitTime, hitTime + 3L, false, hitTime + 1L)))
            .setExternalState("join-state", mapOf(key to stateAfterHit))
            .setWatermarks(defaultWatermarks())
            .build()
    val showResponse = testHarness.doProcess(showRequest)

    val stateAfterShow = showResponse.getExternalState("join-state", key)!!

    // Шаг 3: обработка click-сообщения
    val clickRequest = TestDoProcessRequest.builder("join")
            .setMessages(listOf(buildActionMessage(hitId, hitTime, hitTime + 7L, true, hitTime + 2L)))
            .setExternalState("join-state", mapOf(key to stateAfterShow))
            .setWatermarks(defaultWatermarks())
            .build()
    val clickResponse = testHarness.doProcess(clickRequest)

    val stateAfterClick = clickResponse.getExternalState("join-state", key)!!

    // Шаг 4: срабатывание таймера
    val timer = Timer.builder()
            .setStreamId("timer")
            .setKey(key)
            .setTriggerTimestamp(hitTime + WAIT_SECONDS)
            .build()
    val timerRequest = TestDoProcessRequest.builder("join")
            .setTimers(listOf(timer))
            .setExternalState("join-state", mapOf(key to stateAfterClick))
            .build()
    val timerResponse = testHarness.doProcess(timerRequest)

    // Проверяем финальный результат
    assertEquals(1, timerResponse.getOutputMessagesFlatten().size)
    val result = timerResponse.getOutputMessagesFlatten()[0].getPayload() as JoinedAction
    assertEquals(hitId, result.getHitId())
    assertTrue(result.getClick())

    assertTrue(timerResponse.getExternalState("join-state", key)!!.isReset)
}

Важно

TestComputationHarness не хранит состояние между вызовами doProcess(). Каждый вызов — это независимый батч. Для эмуляции многошаговой обработки необходимо вручную передавать состояние из ответа предыдущего шага в запрос следующего.

Пример: тест со Spring

При использовании Spring Boot Starter тестирование упрощается: PipelineContext создаётся автоматически через FlowAutoConfiguration на основе аннотированных @FlowComputation / @FlowSourceComputation бинов.

Тестовая конфигурация

Для тестов необходимо подменить GrpcServerExecution на NoServerTestExecution, чтобы не запускать реальный gRPC-сервер:

@TestConfiguration
public class MyTestConfiguration {

    @Bean
    public CompanionExecutionConfig companionExecutionConfig() {
        return new CompanionExecutionConfig(0, new MockEnvironmentReader().worker());
    }

    @Bean
    public GrpcServerExecution grpcServerExecution(
            PipelineContext pipelineContext,
            CompanionExecutionConfig companionExecutionConfig
    ) {
        return new NoServerTestExecution(pipelineContext, companionExecutionConfig);
    }
}
@TestConfiguration
class MyTestConfiguration {

    @Bean
    fun companionExecutionConfig(): CompanionExecutionConfig =
        CompanionExecutionConfig(0, MockEnvironmentReader().worker())

    @Bean
    fun grpcServerExecution(
        pipelineContext: PipelineContext,
        companionExecutionConfig: CompanionExecutionConfig
    ): GrpcServerExecution = NoServerTestExecution(pipelineContext, companionExecutionConfig)
}

Ключевые моменты:

  • MockEnvironmentReader — подменяет чтение переменных окружения. Метод .worker() устанавливает YT_FLOW_MODE=Worker.
  • NoServerTestExecution — заглушка для GrpcServerExecution, которая не запускает реальный gRPC-сервер.
  • CompanionExecutionConfig(0, ...) — порт 0 означает, что реальный порт не выделяется.

Тестовый класс

@SpringBootTest(classes = MyTestConfiguration.class)
class JoinFunctionTest {

    @Autowired
    private PipelineContext pipelineContext;

    private TestComputationHarness testHarness;
    private TableSchema keySchema;

    @BeforeEach
    void init() throws IOException {
        this.keySchema = TableSchema.builder()
                .addValue("hash", ColumnValueType.UINT64)
                .addValue("hit_id", ColumnValueType.STRING)
                .addValue("hit_time", ColumnValueType.UINT64)
                .build();

        var specPath = Paths.getSourcePath(
                "yt/yt/flow/examples/java/lb_wait_click_join/test/pipeline.yson");
        var txtSpec = Files.readString(Path.of(specPath));

        this.testHarness = TestComputationHarness.builder()
                .setPipelineContext(pipelineContext)  // инжектированный через Spring
                .setPipelineSpec(txtSpec)
                .build();
    }

    @Test
    void testHitMessageStoresState() {
        var messages = List.of(buildHitMessage("hit-1", 1000L, "payload-data"));

        var request = TestDoProcessRequest.builder("join")
                .setMessages(messages)
                .setWatermarks(Map.of("hit", 0L, "action", 0L))
                .build();

        var response = testHarness.doProcess(request);

        assertTrue(response.getOutputMessagesFlatten().isEmpty());
        assertEquals(1, response.getOutputTimersFlatten().size());
        assertEquals(1, response.getInternalStateSize("join-state"));
    }

    // ... вспомогательные методы и остальные тесты
}
@SpringBootTest(classes = [MyTestConfiguration::class])
class JoinFunctionTest {

    @Autowired
    private lateinit var pipelineContext: PipelineContext

    private lateinit var testHarness: TestComputationHarness
    private lateinit var keySchema: TableSchema

    @BeforeEach
    fun init() {
        keySchema = TableSchema.builder()
                .addValue("hash", ColumnValueType.UINT64)
                .addValue("hit_id", ColumnValueType.STRING)
                .addValue("hit_time", ColumnValueType.UINT64)
                .build()

        val specPath = Paths.getSourcePath(
                "yt/yt/flow/examples/java/lb_wait_click_join/test/pipeline.yson")
        val txtSpec = Files.readString(Path.of(specPath))

        testHarness = TestComputationHarness.builder()
                .setPipelineContext(pipelineContext)  // инжектированный через Spring
                .setPipelineSpec(txtSpec)
                .build()
    }

    @Test
    fun testHitMessageStoresState() {
        val messages = listOf(buildHitMessage("hit-1", 1000L, "payload-data"))

        val request = TestDoProcessRequest.builder("join")
                .setMessages(messages)
                .setWatermarks(mapOf("hit" to 0L, "action" to 0L))
                .build()

        val response = testHarness.doProcess(request)

        assertTrue(response.getOutputMessagesFlatten().isEmpty())
        assertEquals(1, response.getOutputTimersFlatten().size)
        assertEquals(1, response.getInternalStateSize("join-state"))
    }

    // ... вспомогательные методы и остальные тесты
}

Отличия от теста без Spring

Аспект Без Spring Со Spring
Создание PipelineContext Вручную: new PipelineContext() Автоматически через FlowAutoConfiguration
Регистрация Computation pipelineContext.registerComputation(...) Через аннотацию @FlowComputation / @FlowSourceComputation
Регистрация стримов pipelineContext.registerStream(...) Через бин FlowStream<?> или ComputationProvider.getStreams()
Инъекция зависимостей в ProcessFunction Вручную через конструктор Автоматически через @Autowired
Подмена gRPC-сервера Не нужна (нет сервера) NoServerTestExecution + MockEnvironmentReader

Интеграционное тестирование с FlowTestJavaBase

Для полного интеграционного тестирования пайплайна (с реальными C++ воркерами, очередями и стримами) используется базовый класс FlowTestJavaBase.

Зависимости

PEERDIR(
    yt/yt/flow/library/python/integration_test_base
)

Настройка

Тест наследуется от FlowTestJavaBase и задаёт два обязательных атрибута:

from yt.yt.flow.library.python.integration_test_base.yt_flow_java_base import FlowTestJavaBase
import yatest.common

class TestWordCount(FlowTestJavaBase):
    JAVA_RUNNER_BINARY_DIR = yatest.common.binary_path(
        "yt/yt/flow/examples/java/word_count/wordcount/"
    )
    JAVA_MAIN_CLASS = "tech.ytsaurus.flow.examples.wordcount.WordCountApplication"
Атрибут Описание
JAVA_RUNNER_BINARY_DIR Путь к директории с бинарём Java-раннера (содержит run.sh)
JAVA_MAIN_CLASS Полное имя класса точки входа пайплайна: он и запускает пайплайн, и обслуживает компаньон

Примеры интеграционных тестов (Java)

Важно

Интеграционные тесты требуют развёрнутого кластера YTsaurus и запускаются через ya make -ttt. Для быстрой итерации используйте юнит-тесты с TestComputationHarness, описанные выше.

Общие принципы написания интеграционных тестов на пайплайны:

  • Тестируем конечные сценарии. То есть:
    • Записываем входные/исходные данные в локальный YT.
    • В спеке пайплайна отмечаем источники как finite=%true.
    • Для Key-visitor-стрима, который должен работать во время теста, задаём finite=%false.
    • Запускаем пайплайн.
    • Когда визитор должен завершиться, вызываем метод базового класса интеграционного теста self.ask_key_visitor_to_complete("<computation_id>", "<stream_id>"): он переключает динамический параметр finite этого стрима на %true.
    • Ждём, пока пайплайн завершится.
    • Проверяем выходные данные.
  • Готовим окружение тем же кодом, что и в проде:
    • Тестируем тот же бинарь пайплайна что будет работать в проде.
    • Генерируем спеки пайплайна тем же кодом, что будет генерировать их для прода.
    • Если входные/выходные данные нетривиально сериализуются/парсятся, то делаем эту работу общим с продом кодом.
  • Используем общий тестовый фреймворк (есть README.md).
  • Делаем failover тест.
    • Многие ошибки вскрываются на выпадении воркеров и переподхвате их работы другими воркерами. Поэтому делаем тест с problems=True и с более чем 1 воркером.
  • Пишем стабильные тесты.
    • Помним, что в CI любая часть теста может работать неожиданно долго.
      • Идеальное время работы одного теста — 20 секунд при любом типе сборки.
      • Для санитайзерных сборок уменьшаем количество входных данных.
      • Все локальные таймауты выставляем с кратным запасом.
    • Логика выполнения теста не должна значимо зависеть от текущего времени. Например, тест не должен падать, если запустился до полуночи, а завершился — после.

Отладка тестов

Логи

После завершения работы теста его логи можно найти в test-results/py3test/testing_out_stuff в директории теста. Основные важные:

  • run.log — логи python-теста.
  • <test_class_name>/<test_name>/Controller_<number>... — логи контроллеров (.err — это stderr процесса, .log — обычные логи, которые пишутся через YT_LOG_...).
  • <test_class_name>/<test_name>/Worker_<number>... — логи воркеров.
  • <test_class_name>/<test_name>/Runner... — логи раннера.

Если что-то не работает, то имеет смысл поискать ошибки во всех этих логах. Принцип изучения логов такой же, как в рабочем пайплайне; см. логи Vanilla-операции.

Логи можно посмотреть ещё до завершения теста, для этого нужно найти временную директорию, в которой тест работает. Самый простой способ — запускать тест с флагом --keep-temps: ya make --keep-temps -ttt <target>. В этом случае ya make не удалит временную директорию по окончании теста и напечатает ссылку на неё в выводе.

Также можно воспользоваться такими 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)'

# Достать из логов ссылку на UI локального YTsaurus.
alias curlocalyt='cat $(curtestdir)/stderr 2>/dev/null | grep YT'

Тестовый фреймворк

Как настроить окружение для лучшей работы тестового фреймворка и как можно влиять на параметры тестирования можно прочитать в README.md фреймворка.

Поведение интеграционных тестов настраивается через --test-param NAME=VAL:

Параметр

По умолчанию

Значения

Действие

RUNNER_LOG_LEVEL

Error, Info, Debug, …

Уровень логирования процесса, запускающего пайплайн (runner).

PAUSE_BEFORE_FLOW_PROCESS_FEDERATION_TEARDOWN

0

0, 1

Зависнуть в тесте перед остановкой процессов Flow. В комбинации с --test-disable-timeout позволяет надолго оставить работающий локальный YTsaurus и федерацию процессов Flow для неспешного изучения через UI.

EXTERNAL_YT_CONFIG

(не задан)

yson — см. ниже

Запускать пайплайн на реальных внешних кластерах YTsaurus вместо локального рецепта.

Примеры:

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

Локальный YTsaurus всё ещё стартует рецептом, но тест его игнорирует.

Обязательные поля, общие для всех кластеров:

  • path — базовая директория.
  • tablet_cell_bundle — bundle создаваемых динамических таблиц.
  • proxy_role — RPC proxy role.

Опционально: primary_medium (по умолчанию "default").

Список clusters — первый элемент primary. Поля записи: cluster_name (обязательно), proxy_url (по умолчанию равен cluster_name).

Авторизация: при внешнем YTsaurus YT_TOKEN/YT_USER из локального рецепта чистятся, yt-wrapper подхватывает токен из ~/.yt/token.

Изоляция: work_yt_path = path/<local-username>/<test_name>; директория path/<username> удаляется и пересоздаётся один раз на класс в setup_class.

Пример:

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

См. также