Тестирование с 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 необходимо:
PipelineContext— контекст пайплайна с зарегистрированными объектамиComputationи стримами.- Pipeline spec — статическая спецификация пайплайна в формате YSON (файл
pipeline.yson). - 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 секунд при любом типе сборки.
- Для санитайзерных сборок уменьшаем количество входных данных.
- Все локальные таймауты выставляем с кратным запасом.
- Логика выполнения теста не должна значимо зависеть от текущего времени. Например, тест не должен падать, если запустился до полуночи, а завершился — после.
- Помним, что в CI любая часть теста может работать неожиданно долго.
Отладка тестов
Логи
После завершения работы теста его логи можно найти в 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). |
|
|
|
|
|
Зависнуть в тесте перед остановкой процессов Flow. В комбинации с |
|
|
(не задан) |
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>};];}'