Коннекторы YTsaurus Flow
Коннектор — компонент, связывающий пайплайн с объектами YTsaurus (очередью, таблицей и т. д.). Каждый коннектор предоставляет сорс для чтения сообщений и/или синк для записи.
Примечание
Используемые пайплайном коннектеры непосредственно влияют на предоставляемые им гарантии обработки сообщений. Перед выбором коннекторов важно ознакомиться с разделом Гарантии обработки.
Интеграции с внешними (не-YTsaurus) системами описаны в разделе Расширения.
Список коннекторов
|
Коннектор |
Есть сорс |
Есть синк |
Описание |
|
✓ |
✓ |
Чтение и запись в упорядоченную динамическую таблицу с использованием Queue API |
|
|
✓ |
✓ |
Чтение из статических таблиц: фиксированный набор или неограниченная последовательность из директории. Запись непрерывной последовательности статических таблиц в порядке прихода сообщений |
|
|
Random |
✓ |
𐄂 |
Чтение случайных данных, генерируемых на лету. Используется для тестов |
|
✓ |
𐄂 |
Генерация сервис-лога по внешней таблице стейтов. Сервис-лог — это генерация всех ключей из динтаблицы с заданной периодичностью, по сути это способ для переобхода всех стейтов раз в сконфигурированное время |
|
|
𐄂 |
✓ |
Запись в сортированную динамическую таблицу |
Смена источника
Сорс идентифицирует свои партиции по параметрам, описывающим физический источник; конкретный набор таких параметров указан в документации соответствующего коннектора. Эти параметры входят в ключ партиции, поэтому при их изменении сорс описывает уже другой набор партиций.
При смене источника партиции старого источника исчезают из набора и завершаются (Completing ⇒ Completed): связанный с ними стейт, включая запомненные оффсеты, удаляется, а не сохраняется в расчёте на возврат той же партиции (как при интеррапте во время репартиционирования). Партиции нового источника создаются с нуля.
Важно
Параметры, идентифицирующие источник, входят в статический Spec, поэтому менять их можно только на остановленном пайплайне. Новый источник читается независимо от того, что уже было обработано старым: если он содержит те же данные, они будут записаны в выходные потоки повторно.