Работа с ClickHouse в YTsaurus Flow

Запись в ClickHouse реализована в виде необязательного плагина к Flow. Расширение предоставляет только синк, но не источник: данные пишутся напрямую в ClickHouse по нативному протоколу (TCP, порт 9000 по умолчанию), при необходимости через TLS (см. параметр enable_tls). См. исходный код ClickHouse-расширения для Flow.

В основе лежит нативный C++ клиент contrib/libs/clickhouse-cpp. Каждый экземпляр синка держит по одному соединению clickhouse::Client на шард на одной выделенной очереди действий (клиент синхронный и не потокобезопасный). Шарды одного батча пишутся последовательно в этой очереди: базовый упорядоченный синк требует, чтобы батчи завершались строго в порядке их номеров.

Семейство синков

Уровень гарантий задаётся классом синка. TAtMostOnceClickHouseSink сохраняет гарантию at-most-once в обоих режимах стратегии; стратегия выбирает между упорядоченной доставкой с сохранением сообщений и ограниченной очередью с потерями:

  • NYT::NFlow::TClickHouseBatchingSink — рекомендуемый вариант, exactly-once с детерминированным батчингом (один INSERT на батч). Repartition-safe (см. Гарантии обработки).
  • NYT::NFlow::TShardedClickHouseBatchingSink — exactly-once с тем же батчингом для многошардовой формы shard_hosts.
  • NYT::NFlow::TAtLeastOnceClickHouseSink — at-least-once: батч эпохи пишется синхронно в ходе обработки без промежуточного хранения и без токена дедупликации. Ниже задержка и нагрузка на YTsaurus. При ошибке записи эпоха не коммитится, поэтому сообщения не теряются, но при сбое возможны дубликаты. Известные временные ошибки повторяются без ограничения числа попыток, а для неклассифицированных ошибок выполняется не более max_insert_attempts попыток вставки всего (по умолчанию 10). Постоянные ошибки сразу приводят к отказу. Подходит для идемпотентных витрин (ReplacingMergeTree, argMax-роллапы).
  • NYT::NFlow::TAtMostOnceClickHouseSink — at-most-once в обоих режимах. При статическом значении по умолчанию at_most_once_strategy.enabled = false ожидающие сообщения остаются в output_messages и доставляются по порядку, поэтому ошибку подключения или запуска сессии до начала INSERT можно повторить. Значение true включает независимую отправку каждого сообщения через ограниченную очередь в памяти без ожидания доставки. Её динамический лимит задаёт at_most_once_strategy.total_queue_bytes_limit; при отбрасывании сообщений из-за переполнения синк пишет предупреждение. В обоих режимах после начала INSERT синк подтверждает сообщение с ошибкой вместо повторной вставки, поэтому оно не блокирует последующие сообщения.

В обоих exactly-once классах постоянная ошибка вставки или исчерпание попыток для неклассифицированной ошибки приводит к ошибке Giving up insert into ClickHouse. После неё этот экземпляр синка отклоняет все последующие записи, и обработка останавливается. Исправьте целевую таблицу или конфигурацию, затем приостановите и снова запустите пайплайн, чтобы создать новый экземпляр синка.

Синхронного exactly-once синка (TSyncClickHouseSink) не существует: ClickHouse не участвует в табличной транзакции YTsaurus, поэтому exactly-once для ClickHouse обязательно идёт через асинхронный transactional outbox.

Целевая таблица

Exactly-once опирается на блочную дедупликацию ClickHouse по настройке запроса insert_deduplication_token, которую синк передаёт с каждым INSERT. Блочная дедупликация включена по умолчанию только для движков Replicated*MergeTree и SharedMergeTree.

  • Рекомендуется: ReplicatedMergeTree / SharedMergeTree — блочная дедупликация включена по умолчанию.
  • Однохостовый нереплицируемый MergeTree без блочной дедупликации принимается с предупреждением, но exactly-once вырождается в at-least-once. Чтобы включить блочную дедупликацию, задайте non_replicated_deduplication_window.
  • Distributed-таблица как цель отвергается: синк сам маршрутизирует строки по шардам, поэтому целью должна быть локальная таблица каждого шарда. Сервер иначе перенаправил бы блок сам и разорвал связь между токеном дедупликации и шардом.
  • Движки за пределами семейства MergeTree отвергаются.

Каждый экземпляр синка создаёт сессию записи отложенно при первой записи и сохраняет её до своей остановки. При запуске сессии синк подключается ко всем указанным конечным точкам и проверяет их метаданные. Каждая конечная точка должна быть доступна. Внутри шарда у всех конечных точек должны совпадать движок и полная упорядоченная схема (name, type, default_kind, default_expression); схемы разных шардов также должны совпадать. Для Replicated*MergeTree синк проверяет логическую идентичность таблицы: zookeeper_name и zookeeper_path должны совпадать, а replica_name у разных реплик может различаться. Проверка выполняется при любом уровне гарантий доставки.

При динамической перенастройке полная проверка метаданных конечных точек не повторяется. Изменение write_timeout пересоздаёт клиенты и заново проверяет окно дедупликации. При изменении replay_horizon или async_insert обновляется соответствующая проверка окна дедупликации. После изменения целевой таблицы приостановите и снова запустите пайплайн: новые экземпляры синка прочитают и проверят её текущие движок, схему и идентификатор репликации.

Для обычного MergeTree и SharedMergeTree отвергается несколько хостов внутри одного шарда: синк не может доказать, что они предоставляют одну логическую таблицу. Словарь shard_hosts с одним хостом на шард допустим. Однохостовый SharedMergeTree поддерживается. Для однохостового обычного MergeTree выдаётся предупреждение о необходимости настроить non_replicated_deduplication_window.

Окно дедупликации и горизонт реплея

Окно дедупликации конечно. ClickHouse отдельно ограничивает его числом блоков в replicated_deduplication_window и временем в replicated_deduplication_window_seconds; для асинхронных вставок временную границу задаёт replicated_deduplication_window_seconds_for_async_inserts. Значения этих серверных настроек зависят от версии и конфигурации ClickHouse. Реплей, пришедший после вытеснения токена, вставляется повторно, поэтому exactly-once сохраняется, только если реплей доходит до ClickHouse раньше.

Параметр replay_horizon задаёт верхнюю границу лага реплея. При запуске сессии записи синк сравнивает его с серверным replicated_deduplication_window_seconds, а при async_insert = true — с replicated_deduplication_window_seconds_for_async_inserts. Если выбранное окно короче, синк пишет в лог воркера структурированное предупреждение с атрибутами Database, Table, DedupWindowSetting, ReplayHorizon и DedupWindow. Значение окна, заданное для конкретной таблицы, нельзя прочитать через нативный клиент, поэтому при создании таблицы задавайте соответствующее временное окно с запасом относительно replay_horizon.

Шардирование

Один синк умеет писать в несколько независимых шардов. Шард — это набор реплик одной таблицы со своим журналом блочной дедупликации.

  • Однохостовая форма: host + port. Самый простой вариант: одна таблица на одном хосте.
  • Нешардированная форма с несколькими хостами: hosts — плоский список хостов одной таблицы, все с общим port.
  • Многошардовая форма: shard_hosts — словарь «имя шарда → список его хостов». port, database и table общие для всех шардов.

Формы взаимоисключающи: задавать нужно ровно одну из host, hosts и shard_hosts. Для exactly-once нешардированный TClickHouseBatchingSink принимает только host или hosts, а TShardedClickHouseBatchingSink требует shard_hosts. At-least-once и at-most-once синки принимают любую из трёх форм. Список hosts из одного элемента отвергается при валидации спеки — в этом случае используйте host.

"shard_hosts" = {
    "a" = ["ch-a-1"; "ch-a-2"];
    "b" = ["ch-b-1"; "ch-b-2"];
};

Синк маршрутизирует строки сам для корректной организации дедупликации. Он не воспроизводит выражение шардирования какой-либо существующей Distributed-таблицы. Схема корректна для случая «N независимых локальных таблиц, которые читаются через Distributed-таблицу, просто объединяющую их». Если те же таблицы пополняет другой писатель через Distributed, либо читающие запросы полагаются на optimize_skip_unused_shards, локальные JOIN, шардированные словари или distributed_group_by_no_merge, такие запросы будут некорректны.

Выбор хоста внутри шарда

Параметр host_selection_policy задаёт порядок конечных точек внутри каждого шарда. Значение ordered_round_robin используется по умолчанию и сохраняет порядок из спеки. При random_start клиент один раз при создании равновероятно выбирает начальную позицию и циклически сдвигает список. После этого clickhouse-cpp перебирает все конечные точки по кругу; новой случайной перестановки перед каждой попыткой нет.

Политика влияет только на начальную точку failover. Она не меняет маршрутизацию строк, токены дедупликации и отпечаток топологии.

Ключ маршрутизации

sharding_key_columns задаёт колонки, по значениям которых считается ключ маршрутизации. Параметр имеет смысл только вместе с shard_hosts; в нешардированных формах он отвергается при валидации спеки. Если список пуст (по умолчанию), ключом служит идентификатор сообщения: он стабилен при реплее и равномерно распределён, но ко-локации нет — строки с одинаковым бизнес-ключом попадут на разные шарды. Задавайте sharding_key_columns явно, если ко-локация нужна.

Шард выбирается rendezvous-хэшированием по ключу и именам шардов: для каждого шарда считается отпечаток от пары «ключ, имя шарда», и строка уходит на шард с максимальным значением. При добавлении или удалении шарда переезжают только ключи добавленного или удалённого шарда (примерно по числу шардов), остальные остаются на месте.

Токен дедупликации и имена шардов

Токен дедупликации шарда — это токен батча (максимальный идентификатор сообщения в батче) с суффиксом :<имя шарда>. Токен привязан к границе батча и имени шарда, а не к содержимому под-батча, поэтому реплей предъявляет ClickHouse тот же токен с тем же блоком. В нешардированных формах токен остаётся без суффикса.

Отсюда два следствия:

  • Имя шарда — постоянный идентификатор. Переименование шарда так же дорого, как ре-шардирование. Суффикс :<имя шарда> подставляется синком в токен в момент вставки, поэтому все батчи после переименования уходят в ClickHouse уже с новым токеном. Уже вставленные строки при этом никто не переписывает: их токены навсегда остаются в журнале блочной дедупликации ClickHouse под прежним именем и больше ни с одним реплеем не совпадут. Вдобавок имя участвует в маршрутизации, поэтому переименование переносит ключи шарда.
  • Добавление реплики внутрь шарда или замена мёртвого хоста не трогают токены: токен привязан к имени шарда, а не к списку хостов.

Требование к детерминизму полезной нагрузки

Значения колонок sharding_key_columns обязаны быть стабильны при реплее. Если значение ключа у строки изменится, а идентификатор сообщения останется прежним, строка переедет на другой шард: новый шард вставит её (и получится дубликат), а старый получит другой блок под использовавшимся токеном и вероятно отбросит его целиком. Это требование к вышестоящей обработке сообщений; синк его проверить не может.

Защита от смены топологии

Смена набора шардов или sharding_key_columns, пока в состоянии синка остались недоставленные батчи, отвергается с ошибкой при старте. Иначе такие батчи были бы переиграны по новой маршрутизации и часть строк продублировалась бы или потерялась. Защита действует в обеих batching-реализациях: TClickHouseBatchingSink и TShardedClickHouseBatchingSink.

Чтобы перейти с host или hosts на shard_hosts, дождитесь полного слива пайплайна, затем одновременно замените класс на TShardedClickHouseBatchingSink и форму хостов на shard_hosts. Сохраните имя синка и тем самым префикс его состояния. Если после перехода откатить бинарный файл до версии, в реестре которой нет TShardedClickHouseBatchingSink, неизменённая шардированная спека будет отвергнута ещё до чтения состояния batching-синка. Обратная замена класса при наличии недоставленного шардированного состояния не является безопасным откатом.

Ошибка при старте партиции выглядит так:

Refusing to start the ClickHouse sink: the shard topology changed while 3 batch(es) are still undelivered; replaying them under the new topology would duplicate or drop rows. Restore the previous shard_hosts / sharding_key_columns, let the pipeline drain to completion, then apply the change
    persisted_topology_fingerprint = unsharded
    spec_topology_fingerprint      = 3f2a17c9b4e05d81
    oldest_undelivered_batch_bound = 1-4-17

Смотреть её нужно в ошибках партиции во flow view (yt flow get-flow-view или страница пайплайна в UI) и в логе джобы воркера.

Синк хранит в собственном состоянии два вида отпечатков: отпечаток маршрутизации и отпечаток логической цели каждого шарда. При старте они сравниваются с текущей спекой и метаданными ClickHouse. Пока остаются недоставленные батчи, несовпадение любого из них блокирует запуск.

Отпечаток маршрутизации покрывает имена шардов и sharding_key_columns (включая их порядок). Шарды канонически сортируются по имени, поэтому перестановка ключей словаря ничего не меняет. Нешардированные формы (host и hosts) имеют общий константный отпечаток; переход с них на shard_hosts меняет отпечаток и формат токена и блокируется до полного слива.

Списки хостов не входят в отпечаток маршрутизации, но учитываются проверкой логической цели. Для Replicated*MergeTree цель определяется парой zookeeper_name/zookeeper_path: замена мёртвой реплики или изменение списка хостов разрешены, если новые конечные точки указывают на ту же таблицу в Keeper. Для обычного MergeTree и однохостового SharedMergeTree доказать тождество цели по метаданным нельзя, поэтому отпечаток включает базу данных, таблицу, порт и список хостов. Их изменение при недоставленных батчах отвергается: верните прежнюю конечную точку, слейте пайплайн и только затем меняйте конфигурацию. Если прежняя конечная точка потеряна безвозвратно, синк намеренно не переигрывает батчи в другую недоказанную цель автоматически.

Проверка выполняется по партициям: партиция, в которой остались недоставленные батчи, отказывается стартовать, а слитая партиция сразу переходит на новую топологию. «Отказывается» означает не предупреждение, а остановку: Init синка бросает ошибку, джоба партиции падает, и контроллер перезапускает её по обычным правилам обработки падений — то есть партиция падает при каждом старте и не продвигается вовсе, пока спеку не вернут к прежней топологии. Поэтому при обновлении без слива топология может примениться частично, и после возврата к прежней спеке могут отказать уже перешедшие партиции. Меняйте топологию только после полного слива пайплайна (graceful update).

Полный слив защищает недоставленные данные, но не переносит уже записанные строки и не сохраняет историческую ко-локацию ключей. После изменения набора шардов отдельно перераспределите сохранённые данные между шардами, если запросам нужна единая раскладка старых и новых строк.

At-least-once и at-most-once синки не несут токена дедупликации, поэтому описанная защита от смены топологии к ним не применяется.

Маппинг типов

Список колонок, их типы и порядок выводятся автоматически из целевой таблицы: при запуске сессии записи перед первой отправкой синк читает её схему из system.columns. Отправляются только те колонки, которые формирует стрим; колонка таблицы, которой в стриме нет, не отправляется — ClickHouse подставляет её DEFAULT. MATERIALIZED- и ALIAS-колонки пропускаются (их вычисляет ClickHouse). Для каждой отправляемой колонки yson-тип стрима валидируется против типа колонки ClickHouse; несоответствие приводит к падению при запуске сессии записи.

Колонка ClickHouse

Тип yson / YT стрима

Int8

int8

Int16

int16

Int32

int32

Int64

int64

UInt8

uint8

UInt16

uint16

UInt32

uint32

UInt64

uint64

Float32

float

Float64

double

String

string или utf8

FixedString(N)

string или utf8 (длина проверяется при записи)

LowCardinality(String)

string или utf8

Bool

boolean

Date

date

DateTime

datetime

Nullable(T)

опциональный (optional) yson-тип, соответствующий T

Не поддерживаются: Decimal, DateTime64, Enum, UUID, IPv4/IPv6, Array/Tuple/Map и любые композитные / yson-типы. Колонка неподдерживаемого типа допустима в таблице, только если стрим её не пишет (тогда применяется её DEFAULT).

Параметры

Колонки не задаются в спеке — они выводятся автоматически из целевой таблицы.

Ровно одна из форм задания хостов (host, hosts или shard_hosts) должна присутствовать в спеке; остальные параметры подключения общие для всех классов синка.

Параметры статической спеки TClickHouseBatchingSink:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

host

Тип: std::string
Значение по умолчанию: TString("")
Единственный хост ClickHouse в нешардированной форме. Для exactly-once эта форма используется с TClickHouseBatchingSink. Обязателен, если не заданы hosts и shard_hosts: ровно одна из трёх форм задания хостов должна присутствовать в спеке. Хост должен быть доступен и предоставлять таблицу с допустимым движком и полной упорядоченной схемой (name, type, default_kind, default_expression).

port

Тип: unsigned short
Значение по умолчанию: 9000
Нативный TCP-порт. Общий для всех шардов.

hosts

Тип: std::vector<std::string>
Значение по умолчанию: []
Плоский список хостов одной нешардированной таблицы, все с общим port. Для exactly-once эта форма используется с TClickHouseBatchingSink. Обязателен, если не заданы host и shard_hosts. Все хосты должны быть доступны и совпадать по движку и полной упорядоченной схеме (name, type, default_kind, default_expression). Для Replicated*MergeTree должны совпадать zookeeper_name и zookeeper_path, а replica_name может различаться. Несколько хостов с обычным MergeTree или SharedMergeTree отвергаются; однохостовый SharedMergeTree поддерживается. Это правило идентичности целевой таблицы действует для всех гарантий доставки. Список из одного элемента отвергается — используйте host.

host_selection_policy

Тип: NYT::NFlow::EClickHouseHostSelectionPolicy
Значение по умолчанию: ordered_round_robin
Порядок выбора хостов внутри каждого шарда. ordered_round_robin сохраняет порядок из спеки и используется по умолчанию. random_start один раз при создании клиента выбирает равновероятную начальную позицию и циклически сдвигает список. После этого clickhouse-cpp перебирает все конечные точки по кругу. Политика не меняет маршрутизацию строк, токены дедупликации и отпечаток топологии.

shard_hosts

Тип: THashMap<std::string, std::vector<std::string>>
Значение по умолчанию: {}
Словарь «имя шарда → список хостов этого шарда» для многошардовой формы. Обязателен для TShardedClickHouseBatchingSink и отвергается TClickHouseBatchingSink; остальные классы синков принимают любую из трёх форм задания хостов. Все конечные точки должны быть доступны. Внутри каждого шарда должны совпадать движок и полная упорядоченная схема (name, type, default_kind, default_expression); схемы разных шардов также должны совпадать. У реплик должны совпадать zookeeper_name и zookeeper_path, а replica_name может различаться. Несколько хостов с обычным MergeTree или SharedMergeTree отвергаются; однохостовый SharedMergeTree поддерживается. Это правило идентичности целевой таблицы действует для всех гарантий доставки. Имена шардов должны быть непусты и состоять из латинских букв, цифр, _ и -; список хостов шарда не должен быть пуст. Имя шарда — постоянный идентификатор, оно входит в токен дедупликации.

sharding_key_columns

Тип: std::vector<std::string>
Значение по умолчанию: []
Колонки, по значениям которых считается ключ маршрутизации. Требует shard_hosts: в нешардированных формах параметр отвергается при валидации спеки. Пустой список означает маршрутизацию по идентификатору сообщения, то есть без ко-локации.

user

Тип: std::string
Значение по умолчанию: default
Пользователь ClickHouse.

password_env_var

Тип: std::string
Значение по умолчанию: TString("")
Имя переменной окружения с паролем; разрешается через GetEnv при создании клиента. Пустое значение — подключение без пароля.

database

Тип: std::string
Значение по умолчанию: default
База данных целевой таблицы. Общая для всех шардов.

table

Тип: std::string
Обязательный параметр
Целевая таблица. Общая для всех шардов, поэтому <database>.<table> должна существовать на каждом хосте каждого шарда.

codec

Тип: NYT::NFlow::EClickHouseCodec
Значение по умолчанию: lz4
Сжатие нативного протокола.

enable_tls

Тип: bool
Значение по умолчанию: false
Подключаться по TCP+TLS вместо обычного TCP. Порт при этом не меняется: укажите в port защищённый нативный порт сервера (обычно 9440).

tls_ca_files

Тип: std::vector<std::string>
Значение по умолчанию: []
Список путей к файлам корневых сертификатов (CA) для проверки сертификата сервера. Пустой список допустим при выключенном TLS; при включённом TLS он означает использование системных корневых сертификатов. Непустой список требует enable_tls = true; иначе проверка конфигурации завершается ошибкой.

tls_ca_directory

Тип: std::string
Значение по умолчанию: TString("")
Путь к директории с корневыми сертификатами (CA) для проверки сертификата сервера. Пустое значение допустимо при выключенном TLS. Непустой путь требует enable_tls = true; иначе проверка конфигурации завершается ошибкой.

tls_skip_verification

Тип: bool
Значение по умолчанию: false
Пропускать проверку TLS-сессии (сертификат сервера и т. д.). Небезопасно, использовать только для тестовых окружений с самоподписанными сертификатами. Значение true требует enable_tls = true; иначе проверка конфигурации завершается ошибкой. Значение false допустимо при выключенном TLS.

Параметры динамической спеки TClickHouseBatchingSink:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

write_timeout

Тип: TDuration
Значение по умолчанию: 1m
Таймаут на сетевую операцию (send/recv).

retry_backoff

Тип: TDuration
Значение по умолчанию: 1s
Пауза между ретраями вставки.

max_insert_attempts

Тип: long
Значение по умолчанию: 10
Максимум попыток вставки для ошибок, которые нельзя надёжно классифицировать как временные (например, серверные ошибки); после исчерпания запись завершается ошибкой. Известные временные ошибки (сеть, протокол) ретраятся без ограничения, невосстановимые (валидация) не ретраятся.

async_insert

Тип: bool
Значение по умолчанию: false
Включает серверные асинхронные вставки. Все классы синка эмитят async_insert=1 и wait_for_async_insert=1. Batching-синки с гарантией exactly-once дополнительно передают токен дедупликации и эмитят async_insert_deduplicate=1; at-least-once- и at-most-once-синки намеренно не передают токен дедупликации и не задают async_insert_deduplicate.

replay_horizon

Тип: TDuration
Значение по умолчанию: 1d
Верхняя граница лага реплея для проверки окна дедупликации. При запуске сессии записи перед первой вставкой сравнивается с серверным replicated_deduplication_window_seconds либо с replicated_deduplication_window_seconds_for_async_inserts, если включён async_insert. Если выбранное окно короче, в лог воркера пишется предупреждение.

max_rows_per_batch

Тип: long
Значение по умолчанию: 1000
Максимальное число сообщений Flow в батче: batching sink сбрасывает накопленный батч при достижении этого порога. Допустимый диапазон — от 1 до 2000.

max_bytes_per_batch

Тип: long
Значение по умолчанию: 5242880
Максимальный суммарный размер сообщений Flow в батче: batching sink сбрасывает накопленный батч при достижении этого порога. Допустимый диапазон — от 1 байта до 10 MiB.

Параметры статической спеки TShardedClickHouseBatchingSink:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

host

Тип: std::string
Значение по умолчанию: TString("")
Единственный хост ClickHouse в нешардированной форме. Для exactly-once эта форма используется с TClickHouseBatchingSink. Обязателен, если не заданы hosts и shard_hosts: ровно одна из трёх форм задания хостов должна присутствовать в спеке. Хост должен быть доступен и предоставлять таблицу с допустимым движком и полной упорядоченной схемой (name, type, default_kind, default_expression).

port

Тип: unsigned short
Значение по умолчанию: 9000
Нативный TCP-порт. Общий для всех шардов.

hosts

Тип: std::vector<std::string>
Значение по умолчанию: []
Плоский список хостов одной нешардированной таблицы, все с общим port. Для exactly-once эта форма используется с TClickHouseBatchingSink. Обязателен, если не заданы host и shard_hosts. Все хосты должны быть доступны и совпадать по движку и полной упорядоченной схеме (name, type, default_kind, default_expression). Для Replicated*MergeTree должны совпадать zookeeper_name и zookeeper_path, а replica_name может различаться. Несколько хостов с обычным MergeTree или SharedMergeTree отвергаются; однохостовый SharedMergeTree поддерживается. Это правило идентичности целевой таблицы действует для всех гарантий доставки. Список из одного элемента отвергается — используйте host.

host_selection_policy

Тип: NYT::NFlow::EClickHouseHostSelectionPolicy
Значение по умолчанию: ordered_round_robin
Порядок выбора хостов внутри каждого шарда. ordered_round_robin сохраняет порядок из спеки и используется по умолчанию. random_start один раз при создании клиента выбирает равновероятную начальную позицию и циклически сдвигает список. После этого clickhouse-cpp перебирает все конечные точки по кругу. Политика не меняет маршрутизацию строк, токены дедупликации и отпечаток топологии.

shard_hosts

Тип: THashMap<std::string, std::vector<std::string>>
Значение по умолчанию: {}
Словарь «имя шарда → список хостов этого шарда» для многошардовой формы. Обязателен для TShardedClickHouseBatchingSink и отвергается TClickHouseBatchingSink; остальные классы синков принимают любую из трёх форм задания хостов. Все конечные точки должны быть доступны. Внутри каждого шарда должны совпадать движок и полная упорядоченная схема (name, type, default_kind, default_expression); схемы разных шардов также должны совпадать. У реплик должны совпадать zookeeper_name и zookeeper_path, а replica_name может различаться. Несколько хостов с обычным MergeTree или SharedMergeTree отвергаются; однохостовый SharedMergeTree поддерживается. Это правило идентичности целевой таблицы действует для всех гарантий доставки. Имена шардов должны быть непусты и состоять из латинских букв, цифр, _ и -; список хостов шарда не должен быть пуст. Имя шарда — постоянный идентификатор, оно входит в токен дедупликации.

sharding_key_columns

Тип: std::vector<std::string>
Значение по умолчанию: []
Колонки, по значениям которых считается ключ маршрутизации. Требует shard_hosts: в нешардированных формах параметр отвергается при валидации спеки. Пустой список означает маршрутизацию по идентификатору сообщения, то есть без ко-локации.

user

Тип: std::string
Значение по умолчанию: default
Пользователь ClickHouse.

password_env_var

Тип: std::string
Значение по умолчанию: TString("")
Имя переменной окружения с паролем; разрешается через GetEnv при создании клиента. Пустое значение — подключение без пароля.

database

Тип: std::string
Значение по умолчанию: default
База данных целевой таблицы. Общая для всех шардов.

table

Тип: std::string
Обязательный параметр
Целевая таблица. Общая для всех шардов, поэтому <database>.<table> должна существовать на каждом хосте каждого шарда.

codec

Тип: NYT::NFlow::EClickHouseCodec
Значение по умолчанию: lz4
Сжатие нативного протокола.

enable_tls

Тип: bool
Значение по умолчанию: false
Подключаться по TCP+TLS вместо обычного TCP. Порт при этом не меняется: укажите в port защищённый нативный порт сервера (обычно 9440).

tls_ca_files

Тип: std::vector<std::string>
Значение по умолчанию: []
Список путей к файлам корневых сертификатов (CA) для проверки сертификата сервера. Пустой список допустим при выключенном TLS; при включённом TLS он означает использование системных корневых сертификатов. Непустой список требует enable_tls = true; иначе проверка конфигурации завершается ошибкой.

tls_ca_directory

Тип: std::string
Значение по умолчанию: TString("")
Путь к директории с корневыми сертификатами (CA) для проверки сертификата сервера. Пустое значение допустимо при выключенном TLS. Непустой путь требует enable_tls = true; иначе проверка конфигурации завершается ошибкой.

tls_skip_verification

Тип: bool
Значение по умолчанию: false
Пропускать проверку TLS-сессии (сертификат сервера и т. д.). Небезопасно, использовать только для тестовых окружений с самоподписанными сертификатами. Значение true требует enable_tls = true; иначе проверка конфигурации завершается ошибкой. Значение false допустимо при выключенном TLS.

Параметры динамической спеки TShardedClickHouseBatchingSink:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

write_timeout

Тип: TDuration
Значение по умолчанию: 1m
Таймаут на сетевую операцию (send/recv).

retry_backoff

Тип: TDuration
Значение по умолчанию: 1s
Пауза между ретраями вставки.

max_insert_attempts

Тип: long
Значение по умолчанию: 10
Максимум попыток вставки для ошибок, которые нельзя надёжно классифицировать как временные (например, серверные ошибки); после исчерпания запись завершается ошибкой. Известные временные ошибки (сеть, протокол) ретраятся без ограничения, невосстановимые (валидация) не ретраятся.

async_insert

Тип: bool
Значение по умолчанию: false
Включает серверные асинхронные вставки. Все классы синка эмитят async_insert=1 и wait_for_async_insert=1. Batching-синки с гарантией exactly-once дополнительно передают токен дедупликации и эмитят async_insert_deduplicate=1; at-least-once- и at-most-once-синки намеренно не передают токен дедупликации и не задают async_insert_deduplicate.

replay_horizon

Тип: TDuration
Значение по умолчанию: 1d
Верхняя граница лага реплея для проверки окна дедупликации. При запуске сессии записи перед первой вставкой сравнивается с серверным replicated_deduplication_window_seconds либо с replicated_deduplication_window_seconds_for_async_inserts, если включён async_insert. Если выбранное окно короче, в лог воркера пишется предупреждение.

max_rows_per_batch

Тип: long
Значение по умолчанию: 1000
Максимальное число сообщений Flow в батче: batching sink сбрасывает накопленный батч при достижении этого порога. Допустимый диапазон — от 1 до 2000.

max_bytes_per_batch

Тип: long
Значение по умолчанию: 5242880
Максимальный суммарный размер сообщений Flow в батче: batching sink сбрасывает накопленный батч при достижении этого порога. Допустимый диапазон — от 1 байта до 10 MiB.

Параметры статической спеки TAtLeastOnceClickHouseSink:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

host

Тип: std::string
Значение по умолчанию: TString("")
Единственный хост ClickHouse в нешардированной форме. Для exactly-once эта форма используется с TClickHouseBatchingSink. Обязателен, если не заданы hosts и shard_hosts: ровно одна из трёх форм задания хостов должна присутствовать в спеке. Хост должен быть доступен и предоставлять таблицу с допустимым движком и полной упорядоченной схемой (name, type, default_kind, default_expression).

port

Тип: unsigned short
Значение по умолчанию: 9000
Нативный TCP-порт. Общий для всех шардов.

hosts

Тип: std::vector<std::string>
Значение по умолчанию: []
Плоский список хостов одной нешардированной таблицы, все с общим port. Для exactly-once эта форма используется с TClickHouseBatchingSink. Обязателен, если не заданы host и shard_hosts. Все хосты должны быть доступны и совпадать по движку и полной упорядоченной схеме (name, type, default_kind, default_expression). Для Replicated*MergeTree должны совпадать zookeeper_name и zookeeper_path, а replica_name может различаться. Несколько хостов с обычным MergeTree или SharedMergeTree отвергаются; однохостовый SharedMergeTree поддерживается. Это правило идентичности целевой таблицы действует для всех гарантий доставки. Список из одного элемента отвергается — используйте host.

host_selection_policy

Тип: NYT::NFlow::EClickHouseHostSelectionPolicy
Значение по умолчанию: ordered_round_robin
Порядок выбора хостов внутри каждого шарда. ordered_round_robin сохраняет порядок из спеки и используется по умолчанию. random_start один раз при создании клиента выбирает равновероятную начальную позицию и циклически сдвигает список. После этого clickhouse-cpp перебирает все конечные точки по кругу. Политика не меняет маршрутизацию строк, токены дедупликации и отпечаток топологии.

shard_hosts

Тип: THashMap<std::string, std::vector<std::string>>
Значение по умолчанию: {}
Словарь «имя шарда → список хостов этого шарда» для многошардовой формы. Обязателен для TShardedClickHouseBatchingSink и отвергается TClickHouseBatchingSink; остальные классы синков принимают любую из трёх форм задания хостов. Все конечные точки должны быть доступны. Внутри каждого шарда должны совпадать движок и полная упорядоченная схема (name, type, default_kind, default_expression); схемы разных шардов также должны совпадать. У реплик должны совпадать zookeeper_name и zookeeper_path, а replica_name может различаться. Несколько хостов с обычным MergeTree или SharedMergeTree отвергаются; однохостовый SharedMergeTree поддерживается. Это правило идентичности целевой таблицы действует для всех гарантий доставки. Имена шардов должны быть непусты и состоять из латинских букв, цифр, _ и -; список хостов шарда не должен быть пуст. Имя шарда — постоянный идентификатор, оно входит в токен дедупликации.

sharding_key_columns

Тип: std::vector<std::string>
Значение по умолчанию: []
Колонки, по значениям которых считается ключ маршрутизации. Требует shard_hosts: в нешардированных формах параметр отвергается при валидации спеки. Пустой список означает маршрутизацию по идентификатору сообщения, то есть без ко-локации.

user

Тип: std::string
Значение по умолчанию: default
Пользователь ClickHouse.

password_env_var

Тип: std::string
Значение по умолчанию: TString("")
Имя переменной окружения с паролем; разрешается через GetEnv при создании клиента. Пустое значение — подключение без пароля.

database

Тип: std::string
Значение по умолчанию: default
База данных целевой таблицы. Общая для всех шардов.

table

Тип: std::string
Обязательный параметр
Целевая таблица. Общая для всех шардов, поэтому <database>.<table> должна существовать на каждом хосте каждого шарда.

codec

Тип: NYT::NFlow::EClickHouseCodec
Значение по умолчанию: lz4
Сжатие нативного протокола.

enable_tls

Тип: bool
Значение по умолчанию: false
Подключаться по TCP+TLS вместо обычного TCP. Порт при этом не меняется: укажите в port защищённый нативный порт сервера (обычно 9440).

tls_ca_files

Тип: std::vector<std::string>
Значение по умолчанию: []
Список путей к файлам корневых сертификатов (CA) для проверки сертификата сервера. Пустой список допустим при выключенном TLS; при включённом TLS он означает использование системных корневых сертификатов. Непустой список требует enable_tls = true; иначе проверка конфигурации завершается ошибкой.

tls_ca_directory

Тип: std::string
Значение по умолчанию: TString("")
Путь к директории с корневыми сертификатами (CA) для проверки сертификата сервера. Пустое значение допустимо при выключенном TLS. Непустой путь требует enable_tls = true; иначе проверка конфигурации завершается ошибкой.

tls_skip_verification

Тип: bool
Значение по умолчанию: false
Пропускать проверку TLS-сессии (сертификат сервера и т. д.). Небезопасно, использовать только для тестовых окружений с самоподписанными сертификатами. Значение true требует enable_tls = true; иначе проверка конфигурации завершается ошибкой. Значение false допустимо при выключенном TLS.

Параметры динамической спеки TAtLeastOnceClickHouseSink:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

write_timeout

Тип: TDuration
Значение по умолчанию: 1m
Таймаут на сетевую операцию (send/recv).

retry_backoff

Тип: TDuration
Значение по умолчанию: 1s
Пауза между ретраями вставки.

max_insert_attempts

Тип: long
Значение по умолчанию: 10
Максимум попыток вставки для ошибок, которые нельзя надёжно классифицировать как временные (например, серверные ошибки); после исчерпания запись завершается ошибкой. Известные временные ошибки (сеть, протокол) ретраятся без ограничения, невосстановимые (валидация) не ретраятся.

async_insert

Тип: bool
Значение по умолчанию: false
Включает серверные асинхронные вставки. Все классы синка эмитят async_insert=1 и wait_for_async_insert=1. Batching-синки с гарантией exactly-once дополнительно передают токен дедупликации и эмитят async_insert_deduplicate=1; at-least-once- и at-most-once-синки намеренно не передают токен дедупликации и не задают async_insert_deduplicate.

replay_horizon

Тип: TDuration
Значение по умолчанию: 1d
Верхняя граница лага реплея для проверки окна дедупликации. При запуске сессии записи перед первой вставкой сравнивается с серверным replicated_deduplication_window_seconds либо с replicated_deduplication_window_seconds_for_async_inserts, если включён async_insert. Если выбранное окно короче, в лог воркера пишется предупреждение.

Параметры статической спеки TAtMostOnceClickHouseSink:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

host

Тип: std::string
Значение по умолчанию: TString("")
Единственный хост ClickHouse в нешардированной форме. Для exactly-once эта форма используется с TClickHouseBatchingSink. Обязателен, если не заданы hosts и shard_hosts: ровно одна из трёх форм задания хостов должна присутствовать в спеке. Хост должен быть доступен и предоставлять таблицу с допустимым движком и полной упорядоченной схемой (name, type, default_kind, default_expression).

port

Тип: unsigned short
Значение по умолчанию: 9000
Нативный TCP-порт. Общий для всех шардов.

hosts

Тип: std::vector<std::string>
Значение по умолчанию: []
Плоский список хостов одной нешардированной таблицы, все с общим port. Для exactly-once эта форма используется с TClickHouseBatchingSink. Обязателен, если не заданы host и shard_hosts. Все хосты должны быть доступны и совпадать по движку и полной упорядоченной схеме (name, type, default_kind, default_expression). Для Replicated*MergeTree должны совпадать zookeeper_name и zookeeper_path, а replica_name может различаться. Несколько хостов с обычным MergeTree или SharedMergeTree отвергаются; однохостовый SharedMergeTree поддерживается. Это правило идентичности целевой таблицы действует для всех гарантий доставки. Список из одного элемента отвергается — используйте host.

host_selection_policy

Тип: NYT::NFlow::EClickHouseHostSelectionPolicy
Значение по умолчанию: ordered_round_robin
Порядок выбора хостов внутри каждого шарда. ordered_round_robin сохраняет порядок из спеки и используется по умолчанию. random_start один раз при создании клиента выбирает равновероятную начальную позицию и циклически сдвигает список. После этого clickhouse-cpp перебирает все конечные точки по кругу. Политика не меняет маршрутизацию строк, токены дедупликации и отпечаток топологии.

shard_hosts

Тип: THashMap<std::string, std::vector<std::string>>
Значение по умолчанию: {}
Словарь «имя шарда → список хостов этого шарда» для многошардовой формы. Обязателен для TShardedClickHouseBatchingSink и отвергается TClickHouseBatchingSink; остальные классы синков принимают любую из трёх форм задания хостов. Все конечные точки должны быть доступны. Внутри каждого шарда должны совпадать движок и полная упорядоченная схема (name, type, default_kind, default_expression); схемы разных шардов также должны совпадать. У реплик должны совпадать zookeeper_name и zookeeper_path, а replica_name может различаться. Несколько хостов с обычным MergeTree или SharedMergeTree отвергаются; однохостовый SharedMergeTree поддерживается. Это правило идентичности целевой таблицы действует для всех гарантий доставки. Имена шардов должны быть непусты и состоять из латинских букв, цифр, _ и -; список хостов шарда не должен быть пуст. Имя шарда — постоянный идентификатор, оно входит в токен дедупликации.

sharding_key_columns

Тип: std::vector<std::string>
Значение по умолчанию: []
Колонки, по значениям которых считается ключ маршрутизации. Требует shard_hosts: в нешардированных формах параметр отвергается при валидации спеки. Пустой список означает маршрутизацию по идентификатору сообщения, то есть без ко-локации.

user

Тип: std::string
Значение по умолчанию: default
Пользователь ClickHouse.

password_env_var

Тип: std::string
Значение по умолчанию: TString("")
Имя переменной окружения с паролем; разрешается через GetEnv при создании клиента. Пустое значение — подключение без пароля.

database

Тип: std::string
Значение по умолчанию: default
База данных целевой таблицы. Общая для всех шардов.

table

Тип: std::string
Обязательный параметр
Целевая таблица. Общая для всех шардов, поэтому <database>.<table> должна существовать на каждом хосте каждого шарда.

codec

Тип: NYT::NFlow::EClickHouseCodec
Значение по умолчанию: lz4
Сжатие нативного протокола.

enable_tls

Тип: bool
Значение по умолчанию: false
Подключаться по TCP+TLS вместо обычного TCP. Порт при этом не меняется: укажите в port защищённый нативный порт сервера (обычно 9440).

tls_ca_files

Тип: std::vector<std::string>
Значение по умолчанию: []
Список путей к файлам корневых сертификатов (CA) для проверки сертификата сервера. Пустой список допустим при выключенном TLS; при включённом TLS он означает использование системных корневых сертификатов. Непустой список требует enable_tls = true; иначе проверка конфигурации завершается ошибкой.

tls_ca_directory

Тип: std::string
Значение по умолчанию: TString("")
Путь к директории с корневыми сертификатами (CA) для проверки сертификата сервера. Пустое значение допустимо при выключенном TLS. Непустой путь требует enable_tls = true; иначе проверка конфигурации завершается ошибкой.

tls_skip_verification

Тип: bool
Значение по умолчанию: false
Пропускать проверку TLS-сессии (сертификат сервера и т. д.). Небезопасно, использовать только для тестовых окружений с самоподписанными сертификатами. Значение true требует enable_tls = true; иначе проверка конфигурации завершается ошибкой. Значение false допустимо при выключенном TLS.

at_most_once_strategy

Тип: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyParameters>
Необязательная стратегия доставки at-most-once. Поддержка зависит от коннектора; перед включением проверьте документацию выбранного коннектора.

Параметры динамической спеки TAtMostOnceClickHouseSink:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

write_timeout

Тип: TDuration
Значение по умолчанию: 1m
Таймаут на сетевую операцию (send/recv).

retry_backoff

Тип: TDuration
Значение по умолчанию: 1s
Пауза между ретраями вставки.

max_insert_attempts

Тип: long
Значение по умолчанию: 10
Максимум попыток вставки для ошибок, которые нельзя надёжно классифицировать как временные (например, серверные ошибки); после исчерпания запись завершается ошибкой. Известные временные ошибки (сеть, протокол) ретраятся без ограничения, невосстановимые (валидация) не ретраятся.

async_insert

Тип: bool
Значение по умолчанию: false
Включает серверные асинхронные вставки. Все классы синка эмитят async_insert=1 и wait_for_async_insert=1. Batching-синки с гарантией exactly-once дополнительно передают токен дедупликации и эмитят async_insert_deduplicate=1; at-least-once- и at-most-once-синки намеренно не передают токен дедупликации и не задают async_insert_deduplicate.

replay_horizon

Тип: TDuration
Значение по умолчанию: 1d
Верхняя граница лага реплея для проверки окна дедупликации. При запуске сессии записи перед первой вставкой сравнивается с серверным replicated_deduplication_window_seconds либо с replicated_deduplication_window_seconds_for_async_inserts, если включён async_insert. Если выбранное окно короче, в лог воркера пишется предупреждение.

at_most_once_strategy

Тип: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyDynamicParameters>
Значение по умолчанию: {}
Динамические параметры at_most_once_strategy. Поддержка зависит от коннектора; перед настройкой проверьте документацию выбранного коннектора.

Чтобы включить ограниченную очередь at-most-once, задайте at_most_once_strategy.enabled = true в статической спеке. Лимит очереди в байтах настраивается динамически вложенным параметром at_most_once_strategy.total_queue_bytes_limit.

Опция async_insert

async_insert — динамический параметр. При включении все классы синка передают async_insert=1 и wait_for_async_insert=1. Batching-синки с гарантией exactly-once дополнительно передают токен дедупликации и параметр async_insert_deduplicate=1; at-least-once- и at-most-once-синки намеренно не передают токен дедупликации и не задают async_insert_deduplicate.

  • async_insert_deduplicate=1 обязателен для exactly-once (по умолчанию 0 дедупликации нет).
  • wait_for_async_insert=1 позволяет синку узнать об ошибке вставки до продвижения outbox.

Асинхронная дедупликация использует отдельные окна replicated_deduplication_window_for_async_inserts и replicated_deduplication_window_seconds_for_async_inserts. Учтите несовместимость асинхронных вставок с дедупликацией для materialized view.

Пример

В репозитории есть интеграционный тест с настоящим локальным ClickHouse: пайплайн передаёт типизированные строки в заранее созданную таблицу ReplicatedMergeTree и проверяет гарантии каждого класса синка при искусственно вызванных сбоях.

Фрагмент спеки синка:

{
    "sink_class_name" = "NYT::NFlow::TClickHouseBatchingSink";
    "input_stream_ids" = ["rows"];
    "parameters" = {
        "host" = "localhost";
        "port" = 9000;
        "database" = "default";
        "table" = "flow_sink";
    };
}

Шардированный exactly-once синк на двух шардах с ко-локацией по user_id:

{
    "sink_class_name" = "NYT::NFlow::TShardedClickHouseBatchingSink";
    "input_stream_ids" = ["rows"];
    "parameters" = {
        "shard_hosts" = {
            "a" = ["ch-a-1"; "ch-a-2"];
            "b" = ["ch-b-1"; "ch-b-2"];
        };
        "sharding_key_columns" = ["user_id"];
        "port" = 9000;
        "database" = "default";
        "table" = "flow_sink";
    };
}

См. также

Предыдущая