Best Practices
Рекомендации по проектированию, развертыванию и эксплуатации DataFlow Operator для достижения оптимальной производительности, надежности и безопасности.
Проектирование пайплайнов
Выбор между DataFlow и DataFlowCron
Используйте DataFlow когда: - Нужна непрерывная потоковая обработка - Источник — Kafka (streaming) - Требуется минимальная задержка (low latency) - Нет явного "конца" потока данных
Используйте DataFlowCron когда: - Задача должна выполняться по расписанию - Источник — polling база данных (PostgreSQL, ClickHouse, Trino, Nessie) - Нужны post-processing шаги (triggers) - Процесс имеет начало и конец (batch processing)
Архитектура пайплайнов
1. Single Responsibility Каждый DataFlow должен решать одну задачу:
# Хорошо: четкая цель
name: user-events-enrichment
source: kafka://user-events
sink: postgresql://analytics.users
# Плохо: смешение целей
name: everything-pipeline
source: kafka://all-topics # Слишком широко
2. Intermediate Topics Для сложных маршрутов используйте Kafka как промежуточный буфер:
DataFlow A: Source → Kafka Topic A
DataFlow B: Kafka Topic A → Transform → Kafka Topic B
DataFlow C: Kafka Topic B → Sink
3. Error Handling Strategy
# Всегда настраивайте error sink для production
spec:
errors:
type: kafka
config:
brokers: [kafka:9092]
topic: error-messages-${ENV}
ackPolicy: afterWrite
Безопасность
Secrets Management
Никогда не храните credentials в манифестах:
# Плохо: plaintext credentials
sink:
config:
connectionString: "postgres://user:password@host/db" # ❌
# Хорошо: SecretRef
sink:
config:
connectionStringSecretRef:
name: db-credentials
key: connection-string
Создание Secrets:
# Создайте secret из файла
kubectl create secret generic db-credentials \
--from-literal=connection-string="postgres://user:pass@host/db" \
--from-literal=password="secure-password"
# Или из файла
echo -n "secure-password" > password.txt
kubectl create secret generic db-credentials \
--from-file=password=password.txt
rm password.txt
TLS/SSL Configuration
Всегда используйте TLS для production Kafka:
source:
type: kafka
config:
securityProtocol: SASL_SSL
tls:
caFile: /etc/certs/ca.crt
certFile: /etc/certs/client.crt
keyFile: /etc/certs/client.key
sasl:
mechanism: scram-sha-512
usernameSecretRef:
name: kafka-credentials
key: username
passwordSecretRef:
name: kafka-credentials
key: password
Монтирование сертификатов:
apiVersion: v1
kind: Secret
metadata:
name: kafka-certs
type: Opaque
data:
ca.crt: <base64-encoded>
client.crt: <base64-encoded>
client.key: <base64-encoded>
Network Policies
Ограничьте сетевой доступ:
apiVersion: networking.k8s.io/v1
kind: NetworkPolicy
metadata:
name: dataflow-processor
spec:
podSelector:
matchLabels:
app: dataflow-processor
policyTypes:
- Ingress
- Egress
egress:
- to:
- podSelector:
matchLabels:
app: kafka
ports:
- protocol: TCP
port: 9092
- to:
- podSelector:
matchLabels:
app: postgresql
ports:
- protocol: TCP
port: 5432
Производительность
Batch Size Optimization (SQL sinks)
Физический размер write-batch (sink.config.batchSize) и момент source ack (spec.ackGranularity) — разные рычаги. При default collapseBatchOnMessageAck: true значение ackGranularity: message всё ещё принудительно ставит MaxBatchSize = 1 — задайте collapseBatchOnMessageAck: false, если нужны bulk SQL-записи и per-message commit источника. Подробнее: Отказоустойчивость — развязка ack и sink batch.
| Sink | Runtime default (если batchSize не задан) |
Рекомендация | Admission warning если меньше |
|---|---|---|---|
| PostgreSQL | 100 (sinkbatch.DefaultPostgreSQLBatchSize) |
100–500 (до ~1000) | < 100 |
| ClickHouse | 500 (sinkbatch.DefaultClickHouseBatchSize) |
500–1000 | < 500 |
| Trino | 10 (sinkbatch.DefaultTrinoBatchSize) |
≥ 10 (скромнее для крупного JSON / Iceberg) | < 10 |
sink.config — неструктурированный RawExtension, поэтому эти значения подставляет processor, если поле опущено (не OpenAPI default CRD). Единый источник: dataflow/pkg/sinkbatch. batchSize: 0 — flush только по таймеру (batchFlushIntervalSeconds, default 10 через sinkbatch.DefaultBatchFlushIntervalSeconds). Явные значения ниже рекомендованного порога дают admission warning (не reject).
Как пишет каждый SQL sink
PostgreSQL
- Копит до
batchSize(или flush interval), затем одна транзакция. - Путь flush (автоматически):
COPY FROMдля однородного plainINSERT(безON CONFLICT); иначе один multi-rowINSERT … VALUES(включаяON CONFLICT/ upsert); смешанныеUPDATE/гетерогенные формы —pgx.Batch. - Source ack — только после успешного Commit (
AckAfterSuccessfulWrite). - Одно соединение на sink pod (
pgx.Connect, не pool). Предпочтительнее укрупнять batch, чем плодить соединения. - Soft-delete попадает в ту же batch TX (для такого flush —
pgx.Batch).
sink:
type: postgresql
config:
batchSize: 500
batchFlushIntervalSeconds: 10
upsertMode: true
conflictKey: id
ClickHouse
- Default batch 500. Flush: native
PrepareBatch, если доступно; иначе один multi-rowINSERT … VALUESчерез Exec. - Source ack — один раз после успешного flush (не mid-batch). Так сохраняется at-least-once при частичном сбое.
- В raw mode в INSERT идёт
message.Dataas-is (без Unmarshal→Marshal на пути записи). - Крупные batch снижают частоту insert и помогают избежать
TOO_MANY_PARTS. ОбычноbatchFlushIntervalSeconds: 5–10приbatchSize500–1000.
sink:
type: clickhouse
config:
batchSize: 1000
batchFlushIntervalSeconds: 5
upsertMode: true
conflictKey: id
Trino
- Default batch 10. Каждый flush — один HTTP statement (multi-row
INSERT ... VALUES) плюс pollingnextUriдоFINISHED. - Upsert идёт через
MERGE(каталоги Iceberg) — медленнее plain INSERT; admission предупреждает приupsertMode: true. Держите batch скромным и задавайтеqueryTimeoutSecondsна всё окно polling. - Для throughput предпочитайте multi-row INSERT; MERGE — только когда нужны идемпотентные обновления.
sink:
type: trino
config:
batchSize: 10
batchFlushIntervalSeconds: 10
queryTimeoutSeconds: 600
# upsertMode: true # MERGE slow-path — только если нужен
# conflictKey: id
Чеклист тюнинга
- Начните с defaults выше; поднимайте PG/CH batch при устойчивой нагрузке, пока не упрётесь в latency/память.
- Высокий Kafka volume:
ackGranularity: batchилиmessage+collapseBatchOnMessageAck: false. - При ClickHouse
TOO_MANY_PARTSсначала увеличьтеbatchSize/ flush interval. - Для Trino Iceberg/Nessie держите
batchSizeтаким, чтобы один statement укладывался вqueryTimeoutSeconds. - Для Nessie/Iceberg source включайте
incrementalBySnapshotи при больших таблицах —maxRowsPerPoll/maxBytesPerPoll.
Что настраивается в YAML, а что автоматически
| Рычаг | Где | Заметки |
|---|---|---|
channelBufferSize, transformWorkers, ackGranularity, collapseBatchOnMessageAck |
DataFlow spec |
Backpressure и ack vs размер записи |
checkpointSyncOnAck, checkpointSaveInterval |
DataFlow spec |
Политика flush checkpoint; sync-on-ack — async/coalesced (ack path не ждёт ConfigMap Patch) |
Sink batchSize, batchFlushIntervalSeconds |
sink config |
Физический write batch |
Kafka async, compression, flush* |
Kafka sink config |
Defaults: async=true, compression=snappy, flush 100 msgs / 100ms |
incrementalBySnapshot, maxRowsPerPoll, maxBytesPerPoll |
Nessie/Iceberg source config |
Snapshot incremental + лимиты poll |
PG COPY vs multi-VALUES vs pgx.Batch |
Автоматически | По форме flush |
CH PrepareBatch vs Exec |
Автоматически | Native driver при наличии |
| Lakehouse file-delta / parquet fast path | Автоматически | В incremental при added files без deletes |
Buffer Sizing
Каналы между source, transform и sink — буферизованные и блокирующие (по умолчанию 100). Больший буфер сглаживает короткие паузы sink. Смотрите dataflow_channel_fill_ratio{channel="source|routing|…"} при устойчивой насыщенности (CDC и lakehouse семплируют source channel). Flush sink уже перекрывается с накоплением следующего batch (double-buffer). Для CPU-тяжёлых transforms увеличьте transformWorkers (порядок emit сохраняется). Подробности: Архитектура — concurrency пайплайна.
Channel Buffer:
spec:
# Для нормальной нагрузки
channelBufferSize: 100 # default
# Для высоконагруженных потоков
channelBufferSize: 1000
# Для ограниченной памяти
channelBufferSize: 50
Transform workers (параллелизм внутри пода):
spec:
transformWorkers: 4 # по умолчанию 1; диапазон 1–64
replicas > 1 — только для источника Kafka. Для polling/CDC оставляйте replicas: 1 и увеличивайте readBatchSize / batchSize / resources / transformWorkers.
Resource Allocation
Baseline Resources:
spec:
resources:
requests:
cpu: "200m"
memory: "256Mi"
limits:
cpu: "1000m"
memory: "512Mi"
High-Volume Processing:
spec:
resources:
requests:
cpu: "1000m"
memory: "1Gi"
limits:
cpu: "2000m"
memory: "2Gi"
Resource Guidelines: | Scenario | CPU Request | Memory Request | CPU Limit | Memory Limit | |----------|-------------|----------------|-----------|--------------| | Light load | 100m | 128Mi | 500m | 256Mi | | Normal load | 200m | 256Mi | 1000m | 512Mi | | Heavy load | 500m | 512Mi | 2000m | 1Gi | | High volume | 1000m | 1Gi | 2000m | 2Gi |
Polling Source Optimization
PostgreSQL Source:
source:
type: postgresql
config:
# Частый poll для real-time
pollInterval: 5 # seconds
# Редкий poll для batch
pollInterval: 300 # 5 minutes
# Размер читаемого batch
readBatchSize: 1000
# Колонка для отслеживания изменений
changeTrackingColumn: updated_at
orderByColumn: id
Отказоустойчивость
Idempotency Configuration
Всегда настраивайте идемпотентность:
spec:
sink:
type: postgresql
config:
upsertMode: true
conflictKey: id
Different Upsert Strategies:
PostgreSQL:
sink:
type: postgresql
config:
upsertMode: true
conflictKey: id
upsertStrategy: ifNewer # или replace
upsertVersionColumn: updated_at
ClickHouse:
sink:
type: clickhouse
config:
upsertMode: true
conflictKey: id
replacingVersionColumn: updated_at
tableEngine: ReplacingMergeTree
Checkpoint Configuration
Polling Sources (PostgreSQL, ClickHouse, Trino):
spec:
checkpointPersistence: true # default
checkpointSyncOnAck: true # для критичных данных
checkpointSaveInterval: 30s
Kafka Source:
spec:
# checkpointPersistence не нужен для Kafka
# используется consumer group offset
ackGranularity: message # более быстрый commit offset
collapseBatchOnMessageAck: false # сохранить sink.config.batchSize для throughput
Подробнее: Отказоустойчивость — развязка ack и sink batch.
Graceful Shutdown
Termination Grace Period:
spec:
# Достаточно времени для flush batch
terminationGracePeriodSeconds: 600
PreStop Hook:
lifecycle:
preStop:
exec:
command: ["/bin/sh", "-c", "sleep 30"]
Dead Letter Queue Pattern
spec:
errors:
type: kafka
config:
brokers: [kafka:9092]
topic: dlq-${ENV}
ackPolicy: afterWrite
transformations:
- type: filter
config:
condition: "$.required_field" # Проверка обязательных полей
Мониторинг
Key Metrics to Watch
Throughput:
dataflow_processed_messages_total
rate(dataflow_processed_messages_total[5m])
Error Rate:
dataflow_errors_total
rate(dataflow_errors_total[5m])
Lag (для Kafka):
kafka_consumer_group_lag
Latency:
dataflow_processing_duration_seconds
Health Checks
Liveness Probe:
livenessProbe:
httpGet:
path: /health
port: 8080
initialDelaySeconds: 30
periodSeconds: 30
Readiness Probe:
readinessProbe:
httpGet:
path: /ready
port: 8080
initialDelaySeconds: 10
periodSeconds: 10
Alerting Rules
Prometheus Alerts:
groups:
- name: dataflow
rules:
- alert: DataFlowHighErrorRate
expr: rate(dataflow_errors_total[5m]) > 0.1
for: 5m
labels:
severity: warning
annotations:
summary: "High error rate in DataFlow"
- alert: DataFlowNoMessages
expr: rate(dataflow_processed_messages_total[10m]) == 0
for: 10m
labels:
severity: warning
annotations:
summary: "DataFlow stopped processing messages"
- alert: DataFlowKafkaLagHigh
expr: kafka_consumer_group_lag > 10000
for: 5m
labels:
severity: critical
annotations:
summary: "Kafka consumer lag is high"
Эксплуатация
Deployment Strategy
Canary Deployment:
# 1. Создайте новый DataFlow с другим именем
kubectl apply -f dataflow-canary.yaml
# 2. Проверьте метрики
kubectl top pods -l app=dataflow-processor
# 3. Удалите старый и переименуйте канареечный
kubectl delete dataflow old-pipeline
kubectl patch dataflow canary-pipeline -p '{"metadata":{"name":"new-pipeline"}}'
Blue-Green Deployment:
# blue.yaml
metadata:
name: pipeline-blue
labels:
version: blue
# green.yaml
metadata:
name: pipeline-green
labels:
version: green
Backup и Recovery
Backup Checkpoint ConfigMap:
# Экспорт checkpoint
kubectl get configmap df-<name>-checkpoint -o yaml > checkpoint-backup.yaml
# Восстановление
kubectl apply -f checkpoint-backup.yaml
Reset Checkpoint:
# One-shot reset
spec:
checkpointReset: true
Maintenance Windows
Scheduled Maintenance:
spec:
maintenance:
- startTime: "2024-01-15T02:00:00Z"
duration: 2h
repeat: weekly
timezone: UTC
Manual Suspension:
spec:
suspended: true
Тестирование
Unit Testing Transformations
Test JSONPath Expressions:
# Используйте gjson cli или online playground
echo '{"user":{"active":true}}' | gjson "user.active"
# Output: true
Integration Testing
Test DataFlow:
# 1. Создайте DataFlow для теста
kubectl apply -f test-dataflow.yaml
# 2. Отправьте тестовое сообщение
echo '{"test": true}' | kafka-console-producer --topic test-topic
# 3. Проверьте результат
kubectl logs -l app=dataflow-processor --tail=10
# 4. Очистите
kubectl delete dataflow test-dataflow
Load Testing
Generate Load:
# Используйте kcat или kafka-producer-perf-test
kafka-producer-perf-test \
--topic test-topic \
--num-records 1000000 \
--record-size 1000 \
--throughput 10000 \
--producer-props bootstrap.servers=localhost:9092
Monitor Resources:
watch kubectl top pods -l app=dataflow-processor
Checklist для Production
Pre-Deployment
- [ ] Версия Kind правильная (DataFlow vs DataFlowCron)
- [ ] Secrets через
*SecretRef(не plaintext) - [ ] Идемпотентный sink:
upsertMode: true+conflictKey - [ ] Для polling/cron:
checkpointSyncOnAck: true+ upsert - [ ]
replicas: 1для non-Kafka sources - [ ] При CPU-тяжёлых transforms: настроен
transformWorkers(оставьте1, если transforms дешёвые) - [ ] Error sink настроен
- [ ] Resources (requests/limits) установлены
- [ ]
batchSize/ackGranularity/collapseBatchOnMessageAckсбалансированы (message-ack не обязан означатьbatchSize: 1) - [ ] Для Trino:
queryTimeoutSecondsс запасом
Post-Deployment
- [ ] Под в статусе Running
- [ ] Метрики processing/published сообщений
- [ ] Нет ошибок в логах
- [ ] Мониторинг и алерты настроены
- [ ] Alerting rules активны
- [ ] Backup checkpoint настроен
Security
- [ ] TLS для Kafka
- [ ] Network Policies применены
- [ ] RBAC настроен
- [ ] Secrets ротированы
- [ ] Security headers проверены
Documentation
- [ ] Манифест задокументирован
- [ ] Architecture diagram обновлен
- [ ] Runbook создан
- [ ] Contact list для on-call