Тестирование в 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 базовый класс проставляет сам. Для локальной федерации он же прописывает в ресурсы компаньона путь к собранному бинарю, чтобы воркер запустил его с диска.

Пример E2E-теста WordCount

Важно

Интеграционные тесты требуют развёрнутого кластера 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 секунд при любом типе сборки.
      • Для санитайзерных сборок уменьшаем количество входных данных.
      • Все локальные таймауты выставляем с кратным запасом.
    • Логика выполнения теста не должна значимо зависеть от текущего времени. Например, тест не должен падать, если запустился до полуночи, а завершился — после.

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

Логи

После завершения работы теста его логи можно найти в 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>};];}'

См. также

Предыдущая
Следующая