Отказоустойчивость и консистентность данных
DataFlow Operator обрабатывает сообщения с семантикой at-least-once (минимум один раз). При падении или перезапуске пода процессора некоторые сообщения могут быть прочитаны и записаны повторно. В этом документе описано поведение, риски рассинхрона данных и настройка идемпотентных sink для предотвращения дубликатов.
Семантика доставки
- At-least-once: Каждое сообщение доставляется минимум один раз. Дубликаты возможны при перезапуске или падении процессора.
- Exactly-once: Не поддерживается нативно. Используйте идемпотентные sink для достижения effectively-once.
Поведение источников при перезапуске
| Источник | Хранение состояния | При перезапуске |
|---|---|---|
| Kafka | Consumer group (Kafka) | Продолжает с последнего закоммиченного offset. Без дубликатов, если offset закоммичен после записи в sink. |
| PostgreSQL | ConfigMap (по умолчанию); в памяти при checkpointPersistence: false |
По умолчанию продолжает с последней позиции. Без персистенции: перечитывает с начала. |
| PostgreSQL CDC | ConfigMap (lastAckedLSN) |
Продолжает logical replication с последнего acked LSN после записи в sink. |
| ClickHouse | ConfigMap (по умолчанию); в памяти при checkpointPersistence: false |
По умолчанию продолжает с последней позиции. Без персистенции: перечитывает с начала. |
| Trino | ConfigMap (по умолчанию); в памяти при checkpointPersistence: false |
По умолчанию продолжает с последней позиции. Без персистенции: перечитывает с начала. |
| Nessie | ConfigMap при incrementalBySnapshot: true и checkpointPersistence (по умолчанию) |
Инкрементальное чтение по цепочке Iceberg snapshot (file-level delta добавленных data files при возможности); без incrementalBySnapshot — полный scan на каждом poll (checkpoint не используется). |
| Iceberg | ConfigMap при incrementalBySnapshot: true и checkpointPersistence (по умолчанию) |
Как у Nessie; ключ checkpoint — iceberg. |
Горизонтальное масштабирование (spec.replicas)
- Kafka: можно задать
spec.replicas > 1. Все поды используют один consumer group; параллелизм ограничен числом партиций топика. - PostgreSQL, PostgreSQL CDC, ClickHouse, Trino, Nessie:
replicasдолжен быть1(или не задан). Несколько подов с общим checkpoint ConfigMap приведут к дублированию данных. - DataFlowCron:
replicas > 1не поддерживается (один Job процессора на тик расписания).
Внутри одного пода для CPU-тяжёлых transforms увеличивайте spec.transformWorkers (порядок emit сохраняется); sink использует double-buffer flush — см. Архитектура — concurrency пайплайна. Не поднимайте replicas для polling-источников.
Kafka источник
Consumer Kafka помечает offset только после успешной записи сообщения в sink (через msg.Ack()). При ackGranularity: message (см. ниже) offset также сразу коммитится в consumer group.
При падении процессора:
- До записи в sink: Offset не закоммичен. При перезапуске сообщение перечитывается. Дубликата в sink нет.
- После записи в sink, до Ack: Данные могут быть в sink, offset не закоммичен. При перезапуске перечитывание → дубликат в sink.
- После Ack: Offset помечен (и закоммичен при
ackGranularity: message). При перезапуске продолжение со следующего сообщения. Дубликата нет.
Nessie источник (инкрементальный режим)
При source.config.incrementalBySnapshot: true процессор читает только новые Iceberg snapshot с момента последнего Ack. Внутри нового snapshot предпочтительны добавленные data files относительно parent (пустой delta пропускает scan; append-only parquet без deletes — прямой file read). Checkpoint (lastAckedSnapshotID, lastAckedSnapshotSequence) сохраняется в ConfigMap, если включены snapshotCheckpoints (по умолчанию) и spec.checkpointPersistence. Опционально maxRowsPerPoll / maxBytesPerPoll ограничивают объём за poll.
Подробнее: nessie-incremental-snapshots-design.md.
Polling источники (PostgreSQL, ClickHouse, Trino)
Персистенция checkpoint включена по умолчанию. Позиция чтения (lastReadChangeTime, lastReadOrderByValue) сохраняется в ConfigMap df-<name>-checkpoint. После перезапуска источник продолжает с последней позиции после Ack sink. Задайте checkpointPersistence: false, чтобы хранить checkpoint только в памяти (теряется при падении пода).
Legacy-ключи (lastReadID, lastReadTime) мигрируются при загрузке; см. таблицу миграции ниже.
При checkpointPersistence: false при падении пода:
- Состояние теряется.
- При перезапуске источник перечитывает с начала (или с неверной позиции).
- Возможны дубликаты или пропуски в зависимости от момента падения.
Персистенция checkpoint включена по умолчанию. Позиция сохраняется в ConfigMap. При перезапуске источник возобновляет чтение с последней закоммиченной позиции, уменьшая дубликаты. Задайте checkpointPersistence: false в spec, чтобы отключить.
Требуется идемпотентный sink
Для polling источников всегда настраивайте идемпотентный sink (UPSERT, ReplacingMergeTree) для безопасной обработки дубликатов.
Поведение batch sink
PostgreSQL, ClickHouse и Trino sink пишут батчами. Последовательность:
- Накопление сообщений в батч
- Выполнение batch-записи (PostgreSQL оборачивает все statements одной транзакцией и коммитит атомарно)
- Вызов
Ack()для каждого сообщения батча (коммит Kafka offset / продвижение polling checkpoint)
Если процессор падает после успешного commit батча, но до Ack:
- Данные уже в sink
- Позиция источника / checkpoint может быть не продвинута
- При перезапуске: перечитывание → дублирование записей в sink (безопасно с идемпотентным sink)
Уменьшение окна дубликатов
Предпочтительнее меньший sink batchSize при ackGranularity: batch, либо ackGranularity: message.
При message-ack решите, сохранять ли bulk-запись через collapseBatchOnMessageAck.
Гранулярность ack (spec.ackGranularity)
Управляет моментом, когда offset / checkpoint источника считаются зафиксированными относительно записи в sink. Это watermark доставки для источника — не физический размер write-batch sink.
| Значение | Поведение |
|---|---|
batch (по умолчанию) |
После успешного flush batch sink все сообщения батча ack'аются вместе. Kafka source полагается на auto-commit consumer group после MarkMessage. |
message |
После успешной записи каждое сообщение ack'ается отдельно. Для Kafka source после каждого mark вызывается Commit(), чтобы offset consumer group продвигался быстрее. |
Kafka sink всегда ack'ает по сообщению на produce-пути, независимо от ackGranularity.
Длительные INSERT в Trino
Для крупных JSON и таблиц Iceberg/Nessie держите batchSize небольшим (часто 1) и задавайте sink.config.queryTimeoutSeconds так, чтобы покрыть всё время выполнения запроса (включая polling nextUri).
Таймаут на этапе nextUri может произойти уже после старта INSERT в Trino, поэтому повторный запуск может дать дубликаты.
Развязка ack и sink batch (spec.collapseBatchOnMessageAck)
Исторически ackGranularity: message также принудительно заставлял batch sink сбрасывать по одной строке
(MaxBatchSize = 1). Это сужает окно дубликатов, но убивает throughput sink
(N round-trip / N statement вместо одной bulk-записи).
collapseBatchOnMessageAck делает эту связь опциональной:
| Значение | Поведение |
|---|---|
true (по умолчанию) |
Legacy-связь. При ackGranularity: message sink принудительно ставит MaxBatchSize = 1, игнорируя sink.config.batchSize. |
false |
Развязка. Sink сохраняет sink.config.batchSize (и flush interval). После успешного bulk flush процессор всё равно делает per-message source ack (mark / Commit() / checkpoint notify). |
| (игнорируется) | При ackGranularity: batch поле не влияет на размер батча. |
Что контролирует каждый параметр
flowchart LR
subgraph semantic["Семантика доставки"]
AG["ackGranularity\nbatch | message"]
end
subgraph physical["Физический IO sink"]
BS["sink.config.batchSize\n+ flush interval"]
CB["collapseBatchOnMessageAck\ntrue → force size 1"]
end
AG -->|"когда Ack / Commit"| SourceState["Kafka offset /\ncheckpoint progress"]
BS --> Flush["Bulk или single-row write"]
CB -->|"только если ack=message"| Flush
Flush -->|"OnAck после успеха"| AG
| Параметр | Контролирует | Не контролирует |
|---|---|---|
ackGranularity |
Когда фиксируется прогресс источника (на батч vs на сообщение после успешной записи) | Сколько строк уходит в один INSERT/Append |
sink.config.batchSize |
Целевой размер физического write-batch (и часто таймер flush) | Вызывается ли Kafka Commit() на каждый mark |
collapseBatchOnMessageAck |
Переопределяет ли message-ack batchSize в 1 |
Сам debounce checkpoint / факт включения checkpointSyncOnAck |
Матрица режимов
ackGranularity |
collapseBatchOnMessageAck |
Эффективный sink batch | Source ack после flush | Типичный кейс |
|---|---|---|---|---|
batch |
любое | batchSize |
Один проход ack на весь батч | Default / высокий throughput |
message |
true (по умолчанию) |
Принудительно 1 |
На каждое сообщение (и Kafka Commit на mark) |
Минимальное окно re-read; небольшой объём |
message |
false |
batchSize |
На каждое сообщение после успеха bulk | Высокий throughput + более быстрый offset commit |
Последовательность в runtime (message + collapseBatchOnMessageAck: false)
- Sink накапливает до
batchSize(или срабатывает flush interval). - Один bulk
OnFlushпишет батч (транзакция / multi-VALUES / columnar insert — зависит от коннектора). - При успехе
AckAfterSuccessfulWriteпроходит по батчу и ack'ает каждое сообщение. - Для Kafka source: каждый ack →
MarkMessage+Commit()(message granularity). - При
checkpointSyncOnAck: trueкаждый ack может инициировать flush ConfigMap (с coalescing поcheckpointSaveInterval).
Если шаг 2 падает, ни одно сообщение батча не ack'ается → at-least-once перечитывание всего батча при рестарте (нужен идемпотентный sink).
sequenceDiagram
participant Src as Source
participant Sink as Batch sink
participant Ack as Ack / Commit
Src->>Sink: msg 1..N (канал)
Note over Sink: накопление до batchSize
Sink->>Sink: bulk OnFlush
alt flush OK
loop каждое сообщение
Sink->>Ack: Ack(msg i)
Ack->>Src: Mark / Commit / checkpoint notify
end
else ошибка flush
Note over Sink,Ack: без ack — весь батч может replay
end
Окно дубликатов и семантика при crash
Предположим crash после успешной записи в sink, но до завершения всех source ack:
| Режим | Худшее окно re-read (примерно) |
|---|---|
ackGranularity: batch |
До одного полного sink batch |
message + collapse true |
~1 сообщение (write и ack 1:1) |
message + collapse false |
Сообщения из in-flight bulk, уже записанные, но ещё не индивидуально ack'нутые (обычно мало; всё равно ≤ batchSize) |
At-least-once сохраняется во всех режимах: ack никогда не раньше успешной записи в sink.
Идемпотентные sink (upsertMode + conflictKey, ReplacingMergeTree и т.п.) по-прежнему нужны, если дубликаты недопустимы.
Какие sink учитывают collapse / batchSize
collapseBatchOnMessageAck применяется при сборке write loop для batch-ориентированных sink:
- PostgreSQL
- ClickHouse
- Trino
- Nessie / Iceberg
Kafka sink не использует этот batching helper; поведение produce/ack отдельное.
Взаимодействие с checkpointSyncOnAck
- Checkpoint polling / CDC:
Saveкладёт состояние в память;checkpointSyncOnAck: trueвызываетFlushAfterBatchAckпосле ack. Flush неблокирующий (фоновый worker) и coalesced поcheckpointSaveInterval(по умолчанию30s) — ack path не ждёт ConfigMap Get+Patch. - При
ackGranularity: messageи collapsefalseна один bulk flush приходится много ack-колбэков → больше попыток flush, но coalesce + async ограничивают частоту API и не тормозят пайплайн. - Связка высокий объём +
message+checkpointSyncOnAck: trueбез учёта coalesce всё ещё может нагружать API server; для migration/cron чаще достаточноbatchack, либо sync-on-ack с разумнымcheckpointSaveInterval.
Как выбрать режим
| Цель | Рекомендация |
|---|---|
| Обычный streaming | ackGranularity: batch, тюнить batchSize |
| Минимальное окно дубликатов, низкий/средний объём | ackGranularity: message (collapse по умолчанию true) |
| Высокий Kafka → ClickHouse/Postgres throughput и быстрый offset commit | ackGranularity: message + collapseBatchOnMessageAck: false + крупный batchSize + идемпотентный sink |
| Migration / cron polling | Часто checkpointSyncOnAck: true + идемпотентный sink; ackGranularity: batch, если не нужны message-level marks |
Примеры
Legacy message-ack (батч схлопнут до 1):
spec:
ackGranularity: message
# collapseBatchOnMessageAck: true # по умолчанию
sink:
type: postgresql
config:
upsertMode: true
conflictKey: material_id
batchSize: 100 # игнорируется, пока collapse = true
Высокий throughput: message ack без схлопывания sink batch:
spec:
channelBufferSize: 1000
ackGranularity: message
collapseBatchOnMessageAck: false
sink:
type: clickhouse
config:
batchSize: 1000
batchFlushIntervalSeconds: 5
upsertMode: true
conflictKey: event_id
См. также sample dataflow/config/samples/kafka-to-clickhouse-high-volume.yaml и Примеры — high volume.
Настройка идемпотентного sink
PostgreSQL sink
Включите UPSERT, чтобы дубликаты обновляли существующие строки. Batch-записи выполняются в явной транзакции (всё или ничего за flush).
sink:
type: postgresql
config:
connectionString: "postgres://..."
table: output_table
upsertMode: true
conflictKey: id # Опционально; по умолчанию PRIMARY KEY
upsertStrategy: ifNewer # always (по умолчанию) | ifNewer
upsertVersionColumn: updated_at # обязательно при upsertStrategy: ifNewer
Требуется PRIMARY KEY или UNIQUE на колонках конфликта. При upsertStrategy: ifNewer обновление выполняется только если EXCLUDED.<version> > target.<version>.
ClickHouse sink
Включите upsertMode для идемпотентной записи через ReplacingMergeTree (при автосоздании таблицы этот движок используется при upsertMode: true):
sink:
type: clickhouse
config:
connectionString: "clickhouse://..."
table: output_table
upsertMode: true
conflictKey: id
replacingVersionColumn: updated_at
tableEngine: ReplacingMergeTree
Или создайте таблицу вручную:
CREATE TABLE output_table (
id UInt64,
data String,
created_at DateTime DEFAULT now()
) ENGINE = ReplacingMergeTree(created_at)
ORDER BY id;
Дубликаты могут быть видны до background merge; для чтения используйте FINAL или полагайтесь на merge.
Trino sink
Для Iceberg-каталогов включите MERGE-based upsert:
sink:
type: trino
config:
serverURL: "http://trino:8080"
catalog: iceberg # имя каталога должно содержать "iceberg"
schema: default
table: output_table
upsertMode: true
conflictKey: id
При совпадении строка обновляется; ifNewer по колонке версии для Trino пока не поддерживается.
Kafka sink
По умолчанию producer Kafka использует requiredAcks: all, idempotent: true, compression: snappy и async: true (AsyncProducer с flush defaults 100 сообщений / 100ms) для надёжного батчевого produce. Параметры задаются в config sink:
| Поле | По умолчанию | Заметки |
|---|---|---|
requiredAcks |
all |
all | local | none; для idempotent нужен all |
compression |
snappy |
none | gzip | snappy | lz4 | zstd |
idempotent |
true |
При true maxOpenRequests должен быть 1 |
async |
true |
false → SyncProducer (исторический RTT на сообщение) |
flushMessages / flushBytes / flushFrequency |
100 / не задан / 100ms при async | Пороги flush Sarama |
Consumers по-прежнему должны обрабатывать возможные дубликаты (например, идемпотентной обработкой или дедупликацией по ключу) для end-to-end exactly-once.
Рекомендации
- Используйте идемпотентные sink для PostgreSQL (UPSERT), ClickHouse (
upsertMode/ ReplacingMergeTree) и Trino Iceberg (MERGE) при polling источниках или когда возможны дубликаты. - Kafka источник: Consumer group хранит offset; at-least-once сохраняется. Идемпотентный sink для batch sink.
ackGranularity: messageсужает окно re-read; добавьтеcollapseBatchOnMessageAck: false, если нужен крупный sinkbatchSize. - batchSize / ackGranularity / collapseBatchOnMessageAck: Watermark доставки и физический размер записи — разные рычаги; см. Развязка ack и sink batch.
- Migration / cron:
checkpointSyncOnAck: true, идемпотентный sink, при наличии колонки версии —upsertStrategy: ifNewer. - Trino
queryTimeoutSeconds: Таймаут с запасом под пиковую нагрузку. - batchFlushIntervalSeconds: Меньшие интервалы чаще сбрасывают батч.
- Error sink: Настройте
spec.errorsдля неудачных сообщений.
Graceful shutdown
При SIGTERM (например, eviction пода, drain ноды):
- Процессор получает сигнал и отменяет контекст.
- Sink сбрасывают in-flight батчи перед выходом.
PreStop: sleep 5даёт время load balancer перестать направлять трафик.
Убедитесь, что terminationGracePeriodSeconds достаточен для сброса больших батчей (по умолчанию: 600 секунд).
Персистенция checkpoint
По умолчанию включено
Поле checkpointPersistence в spec DataFlow по умолчанию равно true. Явно указывать его не требуется — персистенция checkpoint включена для всех DataFlow с polling-источниками.
Персистенция checkpoint включена по умолчанию. Позиция чтения (lastReadChangeTime, lastReadOrderByValue) сохраняется в ConfigMap df-<name>-checkpoint. При перезапуске процессора polling источники (PostgreSQL, ClickHouse, Trino) возобновляют чтение с последней закоммиченной позиции, уменьшая дубликаты.
Канонический JSON checkpoint для каждого типа источника:
{
"lastReadChangeTime": "2024-06-01T12:00:00.123456789Z",
"lastReadOrderByValue": 5042
}
Legacy-форматы нормализуются при загрузке:
| Legacy | Canonical |
|---|---|
Trino: {"lastReadID": 100} |
{"lastReadOrderByValue": 100} |
ClickHouse: {"lastReadID": 100, "lastReadTime": "..."} |
composite поля выше |
Только time: {"lastReadChangeTime": "..."} |
без изменений (одноколоночный WHERE до появления order key) |
После перезапуска checkpoint Trino/ClickHouse только с lastReadID использует фильтр по order key (WHERE orderByColumn > N) до первого ack с timestamp; затем включается tuple (changeTrackingColumn, orderByColumn) > (time, key).
Чтобы отключить, задайте checkpointPersistence: false:
apiVersion: dataflow.dataflow.io/v1
kind: DataFlow
metadata:
name: my-dataflow
spec:
checkpointPersistence: false # Отключить (по умолчанию: true)
source:
type: postgresql
# ...
Контроллер создаёт ConfigMap и RBAC (ServiceAccount, Role, RoleBinding) для процессора. Checkpoint сохраняется с debounce (по умолчанию раз в 30 секунд) и при graceful shutdown.
Синхронизация checkpoint при ack (spec.checkpointSyncOnAck)
По умолчанию pending checkpoint сбрасывается в ConfigMap по таймеру debounce (checkpointSaveInterval, по умолчанию 30s) и при graceful shutdown. После падения пода polling-источники могут перечитать данные за один интервал debounce.
При checkpointSyncOnAck: true после каждого ack батча sink запрашивается flush checkpoint. Flush асинхронный (не блокирует ack path) и по-прежнему coalesced, не чаще checkpointSaveInterval. Рекомендуется для migration и cron:
spec:
checkpointSyncOnAck: true
checkpointSaveInterval: 5s
source:
type: postgresql
sink:
type: postgresql
config:
upsertMode: true
conflictKey: material_id
Сброс checkpoint
Для повторного полного прогона migration/cron без ручного редактирования df-<name>-checkpoint:
spec:
checkpointReset: true # one-shot; контроллер сбрасывает флаг после reconcile
Или annotation на DataFlow:
metadata:
annotations:
dataflow.dataflow.io/reset-checkpoint: "true"
Процессор очищает checkpoint для типа источника при старте и читает с начала.
Strict idempotency (spec.strictIdempotency)
При strictIdempotency: true admission отклоняет polling-источники с неидемпотентным main sink (без upsertMode). По умолчанию (false) выдаётся warning.
Чеклист
| Сценарий | Рекомендация |
|---|---|
| PostgreSQL sink | upsertMode: true + conflictKey; upsertStrategy: ifNewer при колонке версии |
| ClickHouse sink | upsertMode: true или ручной ReplacingMergeTree + ORDER BY |
| Trino sink (Iceberg) | upsertMode: true + conflictKey |
| Kafka → batch sink | ackGranularity: message (+ опционально collapseBatchOnMessageAck: false для throughput) или меньший batchSize + идемпотентный sink |
| Kafka источник | Идемпотентный sink; ackGranularity: message для быстрого commit offset |
| Polling источники | Идемпотентный sink; checkpointSyncOnAck: true для migration/cron |
| batchSize vs ack | Держите раздельно: см. collapseBatchOnMessageAck — message-ack не всегда означает batchSize: 1 |