Перейти к содержанию

Справочник spec DataFlow

Описание полей DataFlow spec. Оркестрация (Deployment, реконсиляция, status) — в Жизненный цикл и status.

Структура CRD

flowchart TB
  subgraph DataFlow["DataFlow"]
    Spec["spec"]
    Status["status"]
  end

  subgraph SpecFields["поля spec"]
    Source["source (обязательно)"]
    Sink["sink (обязательно)"]
    Trans["transformations (опционально)"]
    Errors["errors (опционально)"]
    Resources["resources (опционально)"]
    Scheduling["scheduling (опционально)"]
    Checkpoint["checkpointPersistence (опционально)"]
    ChannelBuffer["channelBufferSize (опционально)"]
    TransformWorkers["transformWorkers (опционально)"]
    Replicas["replicas (опционально, Kafka)"]
    Image["processorImage / processorVersion (опционально)"]
    Maintenance["maintenance (опционально)"]
  end

  Source --> SourceTypes["type: kafka | postgresql | postgresql-cdc | trino | clickhouse | nessie | iceberg"]
  Sink --> SinkTypes["type: kafka | postgresql | trino | clickhouse | nessie | iceberg"]
  Trans --> TransTypes["timestamp | flatten | filter | mask | router | select | remove | snakeCase | camelCase | debeziumUnwrap | replaceField | headersToPayload | structFlatten | extractField | hoistField | cast | timezone | insertField"]

  Spec --> Source
  Spec --> Sink
  Spec --> Trans
  Spec --> Errors
  Spec --> Resources
  Spec --> Scheduling
  Spec --> Checkpoint
  Spec --> ChannelBuffer
  Spec --> TransformWorkers
  Spec --> Replicas
  Spec --> Image
  Spec --> Maintenance

Поля spec

Поле Обязательность Описание
source Да Тип и конфигурация источника. См. Коннекторы.
sink Да Основной приёмник.
transformations Нет Упорядоченный список трансформаций. См. Трансформации.
errors Нет Error sink для неудачных записей.
resources Нет CPU/память для пода процессора.
nodeSelector, affinity, tolerations Нет Планирование пода.
checkpointPersistence Нет По умолчанию true. Для Nessie — при incrementalBySnapshot: true.
ackGranularity Нет По умолчанию batch. Когда коммитить offset источника относительно успеха sink: batch или message. См. Отказоустойчивость.
collapseBatchOnMessageAck Нет По умолчанию true. При ackGranularity: message — принудительно ли MaxBatchSize = 1. false — сохранить sink.config.batchSize и всё равно ack каждого сообщения после успешного bulk flush. Полное описание: Развязка ack и sink batch.
channelBufferSize Нет По умолчанию 100. Для высокой нагрузки Kafka — 500–1000. См. Архитектура — concurrency пайплайна.
transformWorkers Нет По умолчанию 1. Параллельные transform goroutine в одном поде (1–64). Порядок вывода сохраняется.
replicas Нет По умолчанию 1. > 1 только для Kafka.
processorImage / processorVersion Нет Образ процессора.
imagePullSecrets Нет Pull secrets для пода.
maintenance Нет Окна обслуживания и ручная приостановка процессора. См. ниже.

Окна обслуживания (maintenance)

Секция spec.maintenance задаёт расписание окон обслуживания и ручную приостановку. Пока окно активно или suspended: true, оператор масштабирует Deployment процессора до 0 реплик. Состояние отражается в status.maintenanceStatus — см. Жизненный цикл и status.

Поле Обязательность Описание
startTime При расписании Начало окна в формате RFC3339 (напр. 2024-01-01T02:00:00Z).
duration При расписании Длительность окна в формате Go duration (напр. 2h, 30m).
repeat Нет Повтор: daily, weekly, monthly. Пустое значение — одноразовое окно.
timezone Нет IANA timezone (напр. Europe/Moscow). По умолчанию UTC.
description Нет Текстовое описание окна.
suspended Нет true — ручная остановка процессора (аналог кнопки «Остановить» в GUI).

Поля startTime и duration задаются вместе. Validating webhook проверяет формат RFC3339, корректность duration и timezone.

spec:
  maintenance:
    startTime: "2024-01-01T02:00:00Z"
    duration: "2h"
    repeat: daily
    timezone: Europe/Moscow
    description: "Ночное обслуживание БД"

Ручная приостановка без расписания:

spec:
  maintenance:
    suspended: true

Управление через Web GUI: кнопки Stop/Start для отдельного потока, Stop all / Start all для namespace.

Секреты

Credentials через SecretRef — см. Коннекторы — Secrets.

Валидация

При включённом validating webhook невалидный spec отклоняется на admission.

Те же правила применяются к встроенному DataFlowSpec в DataFlowCron.

См. также