Запись в 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";
};
}
См. также