Тестирование в YTsaurus Flow (Go)
Примечание
Данная страница описывает юнит-тестирование компьютейшенов Go-пайплайна через харнесс flowtest, а также интеграционное тестирование полного пайплайна через FlowTestGoBase.
Общая архитектура тестирования
В продакшене C++ воркер отправляет gRPC-запросы компаньону, передавая сообщения, таймеры, визиты, стейты и вотермарки. Компаньон разбирает запрос, собирает по нему flow.Job и flow.Runtime и вызывает Process Function зарегистрированного компьютейшена.
В юнит-тестах место воркера занимает flowtest.Harness из пакета flowtest. Харнесс хранит то, что воркер сообщает компаньону, — стримы, схему ключа, объявленные стейты и параметры — и прогоняет компьютейшен через ту же джобу, тот же рантайм и ту же диспетчеризацию, что и сервер компаньона, вплоть до рендеринга ответа в проволочный формат. Поэтому сообщение в необъявленный стрим и незакодируемый ключ падают уже в юнит-тесте, а не в джобе.
Тестируется то же значение *flow.Computation, которое регистрируется в пайплайне: харнессу передаётся результат flow.NewRowComputation или родственного конструктора, а сорс отличается от трансформа только тем, чем он был создан, — сообщать об этом харнессу отдельно не нужно.
Ни кластера, ни gRPC-соединения, ни flow_server для юнит-тестов не требуется. Тесты пишутся на стандартном testing; в примерах для проверок используется testify (require).
Зависимости
Отдельный PEERDIR для харнесса не нужен: зависимости Go-модуля выводятся из импортов. Достаточно перечислить тестовые файлы в GO_TEST_SRCS модуля пайплайна:
GO_PROGRAM()
SUBSCRIBER(
g:yt-flow
)
SRCS(
main.go
word_count_mapper.go
)
GO_TEST_SRCS(
word_count_mapper_test.go
)
END()
RECURSE_FOR_TESTS(
gotest
test
)
И добавить рядом директорию gotest с модулем GO_TEST_FOR, через который тесты запускаются:
GO_TEST_FOR(yt/yt/flow/examples/go/word_count)
SUBSCRIBER(
g:yt-flow
)
SIZE(SMALL)
END()
Тестирование Process Function
Создание харнесса
Харнесс создаётся функцией flowtest.New(tb, computation, opts). Первый аргумент — *testing.T (подойдёт также *testing.B и *testing.F): обо всех ошибках использования харнесс сообщает через него, поэтому в тесте остаётся только то, что тест утверждает.
h := flowtest.New(t, flow.NewRowComputation("mapper", &wordCountMapper{}), flowtest.Options{
Streams: map[string]flow.Schema{"words": flowtest.Schema("word:string")},
KeySchema: flowtest.Schema("word:string"),
InternalStates: []string{wordStateName},
})
Поля flowtest.Options:
| Поле | Описание |
|---|---|
Streams |
Стримы, по которым компьютейшен обменивается сообщениями, по идентификатору стрима. Читать и писать компьютейшен может только перечисленные здесь стримы. |
KeySchema |
Схема ключа, по которому сгруппированы входы. Компьютейшен без группировки оставляет поле пустым. |
InternalStates |
Имена внутренних стейтов, которые объявляет компьютейшен, — они доезжают до него как parameters.internal_states. |
ExternalStates |
Схемы внешних стейтов, которыми компьютейшен владеет, по имени стейта. Имена — абсолютные пути, как того требует воркер. |
JoinedExternalStates |
Схемы внешних стейтов, которые компьютейшен читает, не владея ими. |
Parameters |
Карта parameters статической спеки — то, что компьютейшен читает через rt.Parameters(). |
DynamicParameters |
Карта parameters динамической спеки. |
Схема колонок собирается хелпером flowtest.Schema("word:string", "count:int64") — имена типов те же, что в YTsaurus. Схему, которую так не описать, стройте через flow.NewSchema из schema.Schema.
Для типизированного YSON-стрима используйте ту же схему, которую регистрирует пайплайн: flow.YSONMessageSchema[event](). Так структура со встроенным flow.YSONMessage описывает и колонки спеки, и входные строки теста.
Для обычной структуры без flow.YSONMessage остаётся flowtest.SchemaOf(event{}). Он следует общему schema.Infer: в частности, Go-строка становится utf8. Если схема должна дословно совпасть с уже существующей спекой, используйте flowtest.Schema.
Входы
Входы одного батча строятся методами харнесса и передаются в Process одним вызовом:
| Метод | Что строит |
|---|---|
h.Key(flowtest.Row{...}) |
Ключ по схеме KeySchema. |
h.Message(streamID, row) |
Сообщение без ключа — то, что получает компьютейшен без группировки. |
h.KeyedMessage(streamID, key, row) |
Сообщение вместе с ключом, по которому оно сгруппировано. |
h.Timer(key, triggerTimestamp) |
Сработавший таймер ключа. |
h.Visit(key) |
Визит ключа из key-visitor стрима. |
h.SetWatermark(streamID, watermark) |
Вотермарк стрима; держится до следующей установки. |
Каждому сообщению выдаётся собственный идентификатор, как это делает воркер. Таймстемпы остаются нулевыми: тесту, которому они нужны, достаточно проставить их на результате.
msg := h.KeyedMessage("hits", key, flowtest.Row{"hit_id": "h1"})
msg.EventTimestamp = 1000
h.Process(inputs ...flow.Input) прогоняет компьютейшен по батчу и возвращает *flowtest.Response; если обработка вернула ошибку, тест падает. Стейт переживает прогон: то, что компьютейшен записал, применяется к стейту следующего прогона — ровно так воркер применяет дельту ответа перед отправкой очередного батча. Тесту, которому нужен чистый лист, следует построить новый харнесс.
Полный пример
Юнит-тесты маппера из WordCount — харнесс, батч сообщений и проверка внутреннего стейта:
func newHarness(t *testing.T) *flowtest.Harness {
return flowtest.New(t, flow.NewRowComputation("mapper", &wordCountMapper{}), flowtest.Options{
Streams: map[string]flow.Schema{"words": flowtest.Schema("word:string")},
KeySchema: flowtest.Schema("word:string"),
InternalStates: []string{wordStateName},
})
}
func TestRepeatedWordAccumulates(t *testing.T) {
h := newHarness(t)
key := h.Key(flowtest.Row{"word": "hello"})
var batch []flow.Input
for range 3 {
batch = append(batch, h.KeyedMessage("words", key, flowtest.Row{"word": "hello"}))
}
r := h.Process(batch...)
require.EqualValues(t, 3, counterOf(t, r, key).Count)
}
func TestCounterSurvivesTheBatch(t *testing.T) {
h := newHarness(t)
key := h.Key(flowtest.Row{"word": "hello"})
h.Process(h.KeyedMessage("words", key, flowtest.Row{"word": "hello"}))
r := h.Process(h.KeyedMessage("words", key, flowtest.Row{"word": "hello"}))
require.EqualValues(t, 2, counterOf(t, r, key).Count)
}
func TestWordsAreCountedApart(t *testing.T) {
h := newHarness(t)
hello := h.Key(flowtest.Row{"word": "hello"})
world := h.Key(flowtest.Row{"word": "world"})
r := h.Process(
h.KeyedMessage("words", hello, flowtest.Row{"word": "hello"}),
h.KeyedMessage("words", world, flowtest.Row{"word": "world"}),
h.KeyedMessage("words", hello, flowtest.Row{"word": "hello"}),
)
require.EqualValues(t, 2, counterOf(t, r, hello).Count)
require.EqualValues(t, 1, counterOf(t, r, world).Count)
}
func counterOf(t *testing.T, r *flowtest.Response, key flow.Payload) wordCountState {
t.Helper()
var counter wordCountState
require.True(t, r.InternalStateYSON(wordStateName, key, &counter), "no counter stored for the key")
return counter
}
Ошибки обработки
Ошибка, возвращённая обработчиком, прекращает обработку батча целиком: воркер повторит запрос, поэтому частичного ответа не бывает. Такой прогон проверяется методом h.ProcessError, который возвращает ошибку и падает, если обработка, наоборот, прошла успешно:
err := h.ProcessError(h.Message("queue", flowtest.Row{"data": "}not json{"}))
require.ErrorContains(t, err, "parsing the data column")
Прогон, завершившийся ошибкой, не производит вывода и не меняет стейт — поэтому ProcessError и не возвращает Response.
Таймеры и вотермарки
Таймер строится по ключу и времени срабатывания. Пайплайну с несколькими таймерными стримами нужный выбирается полем StreamID — пустое значение означает единственный таймерный стрим пайплайна:
timer := h.Timer(key, closeTime)
timer.StreamID = timerStream
r := h.Process(timer)
Вотермарк стрима задаётся h.SetWatermark и держится для всех последующих прогонов. Так проверяется отбрасывание опоздавших данных: компьютейшен читает rt.MinWatermark(), минимум по входным стримам, поэтому не продвинувшийся стрим удерживает окно открытым для остальных.
h.SetWatermark(hitStream, hitTime+3)
h.SetWatermark(actionStream, 0)
Полный набор тестов на окно с таймером и вотермарками — в Wait Click Join.
Тестирование стейтов
Стейт, с которым компьютейшен начинает прогон, кладётся в харнесс до вызова Process, а результат читается из Response. Подробнее о самих аксессорах — в разделе State Accessor.
Internal state
Имя внутреннего стейта должно быть объявлено в InternalStates, иначе харнесс сообщит ровно ту же ошибку, что и рантайм в джобе.
| Метод | Что кладёт |
|---|---|
h.PutInternalState(name, key, data) |
Сырые байты, которые читает flow.OpenRawState. |
h.PutInternalStateYSON(name, key, value) |
Значение, сериализованное в YSON, — то, что читает flow.OpenYSONState. |
h.PutInternalStateProto(name, key, value) |
Сериализованное protobuf-сообщение для flow.OpenProtoState. |
Обратно стейт читается методами Response:
var counter wordCountState
require.True(t, r.InternalStateYSON(wordStateName, key, &counter))
require.EqualValues(t, 1, counter.Count)
External state
Внешний стейт, которым компьютейшен владеет, кладётся h.PutExternalState(name, key, row), а читается как строка — r.ExternalState возвращает flow.Payload, r.ExternalStateRow — уже декодированную flowtest.Row.
Примечание
Внутренний стейт и присоединённый внешний стейт доезжают до прогона только для тех ключей, для которых что-то хранят. Внешний стейт, которым компьютейшен владеет, доезжает для каждого ключа батча, пустым там, где не хранится ничего: воркер резолвит собственный стейт для каждого переданного ключа — именно это позволяет писать стейт для ключа, увиденного впервые.
Юнит-тесты редьюсера из Shuffle, который считает события во внешнем стейте:
var shuffleStreams = []string{"event_a", "event_b", "event_c", "event_d"}
func newReducerHarness(t *testing.T) *flowtest.Harness {
streams := make(map[string]flow.Schema, len(shuffleStreams))
for _, streamID := range shuffleStreams {
streams[streamID] = eventSchema
}
return flowtest.New(t, flow.NewRowComputation("reducer", &eventReducer{}), flowtest.Options{
Streams: streams,
KeySchema: flowtest.Schema("value:string"),
ExternalStates: map[string]flow.Schema{shuffleStateName: flowtest.Schema("count:int64")},
})
}
func TestAValueIsCountedOncePerShuffleStream(t *testing.T) {
h := newReducerHarness(t)
key := h.Key(flowtest.Row{"value": "v"})
var batch []flow.Input
for _, streamID := range shuffleStreams {
batch = append(batch, h.KeyedMessage(streamID, key, flowtest.Row{"value": "v"}))
}
r := h.Process(batch...)
require.EqualValues(t, 4, countOf(t, r, key))
require.Empty(t, r.Messages())
require.Empty(t, r.Timers())
}
func TestValuesAreCountedApart(t *testing.T) {
h := newReducerHarness(t)
first := h.Key(flowtest.Row{"value": "v1"})
second := h.Key(flowtest.Row{"value": "v2"})
r := h.Process(
h.KeyedMessage("event_a", first, flowtest.Row{"value": "v1"}),
h.KeyedMessage("event_b", second, flowtest.Row{"value": "v2"}),
h.KeyedMessage("event_c", first, flowtest.Row{"value": "v1"}),
)
require.EqualValues(t, 2, countOf(t, r, first))
require.EqualValues(t, 1, countOf(t, r, second))
}
func TestCounterSurvivesTheBatch(t *testing.T) {
h := newReducerHarness(t)
key := h.Key(flowtest.Row{"value": "v"})
h.Process(h.KeyedMessage("event_a", key, flowtest.Row{"value": "v"}))
r := h.Process(h.KeyedMessage("event_b", key, flowtest.Row{"value": "v"}))
require.EqualValues(t, 2, countOf(t, r, key))
}
func countOf(t *testing.T, r *flowtest.Response, key flow.Payload) int64 {
t.Helper()
row, ok := r.ExternalState(shuffleStateName, key)
require.True(t, ok, "no counter stored for the key")
count, err := row.Int64(countColumn)
require.NoError(t, err)
return count
}
Joined external state
Присоединённый внешний стейт — стейт, который компьютейшен читает, не владея им, — кладётся h.PutJoinedExternalState(name, key, row) и читается через r.JoinedExternalState / r.JoinedExternalStateRow. Записать в него нельзя: ничего записанного в read-only стейт из ответа не выходит.
h := flowtest.New(t, flow.NewRowComputation("lookup_join", &lookupJoin{}), flowtest.Options{
Streams: map[string]flow.Schema{"event": flowtest.Schema("key:uint64")},
KeySchema: flowtest.Schema("hash:uint64", "key:uint64"),
JoinedExternalStates: map[string]flow.Schema{
referenceStateName: flowtest.Schema("hash:uint64", "key:uint64", "name:string"),
},
})
h.PutJoinedExternalState(referenceStateName, key, flowtest.Row{"key": uint64(1), "name": "alice"})
Ключ, для которого строка не положена, до компьютейшена не доезжает — как и в проде, где воркер джойнит то, что нашёл, и ничего сверх того. Полный набор тестов — в external_state_join.
Анализ результатов
*flowtest.Response, возвращённый Process, — это то, что произвёл прогон: собранный вывод и стейты в том виде, в каком они будут сохранены.
| Метод | Что возвращает |
|---|---|
Groups() |
[]flow.OutputGroup — группы вывода в порядке появления. |
Messages() |
Выходные сообщения всех групп, по порядку. |
MessagesOn(streamID) |
Выходные сообщения одного стрима. |
Rows() |
Пейлоады выходных сообщений, декодированные в flowtest.Row и выровненные с Messages(). |
Distribute() |
Флаг distribute каждого сообщения, выровненный с Messages(). |
Timers() |
[]flow.TimerRequest — таймеры, которые компьютейшен попросил воркер поставить. |
Группа вывода — это происхождение (lineage) вывода, а не форма входа: RowFunction открывает по группе на вход, BatchFunction — одну на батч, а группы, в которые ничего не записали, отбрасываются.
Стейты читаются так:
| Метод | Что возвращает |
|---|---|
InternalStateRaw(name, key) |
Байты, которые внутренний стейт хранит для ключа. |
InternalStateYSON(name, key, dst) |
Десериализует в dst YSON внутреннего стейта. |
InternalStateProto(name, key, dst) |
Десериализует в dst protobuf-сообщение внутреннего стейта. |
InternalStateReset(name, key) |
Прогон очистил стейт ключа. |
InternalStateWritten(name) |
Прогон писал в стейт: до воркера доезжает только записанное. |
InternalStateLen(name) |
Число ключей, для которых стейт читался или писался. |
ExternalState(name, key), ExternalStateRow(name, key) |
Строка внешнего стейта — как flow.Payload и как flowtest.Row. |
ExternalStateReset(name, key), ExternalStateWritten(name), ExternalStateLen(name) |
То же для внешнего стейта. |
JoinedExternalState(name, key), JoinedExternalStateRow(name, key) |
Строка присоединённого внешнего стейта. |
Стейт рапортуется таким, каким он будет сохранён: очищенная прогоном запись читается как отсутствующая, а отличить её от той, которой никогда не было, позволяет *Reset.
Запуск юнит-тестов
Юнит-тесты — это тесты размера SMALL, кластер им не нужен.
cd yt/yt/flow/examples/go/word_count
go test ./...
Отфильтровать один тест можно по имени:
go test ./... -run 'TestCounterSurvivesTheBatch'
Интеграционное тестирование с FlowTestGoBase
Для полного интеграционного тестирования пайплайна (с реальными C++ воркерами, очередями и стримами) используется базовый класс FlowTestGoBase — Python-тест, запускающий тот же Go-бинарь, что поедет в прод.
Запускает пайплайн в таком тесте раннер, а не сам тест: Go-бинарь стартует как ./word_count --config pipeline.yson --flow-bin flow_server, обогащает спеку и передаёт управление flow_server, который её устанавливает. Компаньон в джобе поднимает воркер — ровно как в проде.
Зависимости
Интеграционному тесту нужны рецепт кластера, DEPENDS на бинарь пайплайна и flow_server, а также DATA со спекой. Полный ya.make теста из WordCount:
PY3TEST()
INCLUDE(${ARCADIA_ROOT}/yt/yt/flow/library/python/integration_test_base/recipe.inc)
TEST_SRCS(
test_wordcount.py
yt_sync.py
)
PEERDIR(
yt/yt/flow/library/python/queue
)
DEPENDS(
${MODDIR}/..
yt/yt/flow/bin/flow_server
)
DATA(arcadia/${MODDIR}/pipeline.yson)
REQUIREMENTS(
cpu:4
ram:32
)
TAG(ya:huge_logs)
SIZE(MEDIUM)
END()
Настройка
Тест наследуется от FlowTestGoBase и задаёт атрибут GO_COMPANION_BINARY:
class Test(FlowTestGoBase):
GO_COMPANION_BINARY = yatest.common.binary_path("yt/yt/flow/examples/go/word_count/word_count")
| Атрибут | Описание |
|---|---|
GO_COMPANION_BINARY |
Путь к бинарю Go-пайплайна: он же раннер, он же компаньон. |
VANILLA_WORKER_PORT_COUNT |
Число портов на воркер; по умолчанию 3 — rpc, мониторинг и порт, на котором воркер поднимает компаньон. |
Пайплайн запускается методом start_flow_process_federation, которому спека передаётся аргументом --config; --flow-bin базовый класс проставляет сам. Для локальной федерации он же прописывает в ресурсы компаньона путь к собранному бинарю, чтобы воркер запустил его с диска.
Важно
Интеграционные тесты требуют развёрнутого кластера YTsaurus и относятся к размеру MEDIUM, поэтому запускаются через ya test -tt. Для быстрой итерации используйте юнит-тесты, описанные выше.
Общие принципы написания интеграционных тестов на пайплайны:
- Тестируем конечные сценарии. То есть:
- Записываем входные/исходные данные в локальный 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>};];}'
См. также
- Computation (Go)
- Работа со стейтами (Go)
- State Accessor (Go)
- Примеры: Word Count (Go)
- Если дорабатываете сам Flow — Фреймворк для тестирования пайплайнов.