Коннекторы YTsaurus Flow

Коннектор — компонент, связывающий пайплайн с объектами YTsaurus (очередью, таблицей и т. д.). Каждый коннектор предоставляет сорс для чтения сообщений и/или синк для записи.

Примечание

Используемые пайплайном коннектеры непосредственно влияют на предоставляемые им гарантии обработки сообщений. Перед выбором коннекторов важно ознакомиться с разделом Гарантии обработки.

Интеграции с внешними (не-YTsaurus) системами описаны в разделе Расширения.

Список коннекторов

Коннектор

Есть сорс

Есть синк

Описание

Queue

✓

✓

Чтение и запись в упорядоченную динамическую таблицу с использованием Queue API

Static Table

✓

✓

Чтение из статических таблиц: фиксированный набор или неограниченная последовательность из директории. Запись непрерывной последовательности статических таблиц в порядке прихода сообщений

Random

✓

𐄂

Чтение случайных данных, генерируемых на лету. Используется для тестов

Сервис-лог

✓

𐄂

Генерация сервис-лога по внешней таблице стейтов. Сервис-лог — это генерация всех ключей из динтаблицы с заданной периодичностью, по сути это способ для переобхода всех стейтов раз в сконфигурированное время

Sorted Dynamic Table

𐄂

✓

Запись в сортированную динамическую таблицу

Смена источника

Сорс идентифицирует свои партиции по параметрам, описывающим физический источник; конкретный набор таких параметров указан в документации соответствующего коннектора. Эти параметры входят в ключ партиции, поэтому при их изменении сорс описывает уже другой набор партиций.

При смене источника партиции старого источника исчезают из набора и завершаются (Completing ⇒ Completed): связанный с ними стейт, включая запомненные оффсеты, удаляется, а не сохраняется в расчёте на возврат той же партиции (как при интеррапте во время репартиционирования). Партиции нового источника создаются с нуля.

Важно

Параметры, идентифицирующие источник, входят в статический Spec, поэтому менять их можно только на остановленном пайплайне. Новый источник читается независимо от того, что уже было обработано старым: если он содержит те же данные, они будут записаны в выходные потоки повторно.