Getting Started
Установка оператора, первый потоковый или scheduled-конвейер, либо локальная среда разработки.
Когда читать: первая установка или первый pipeline. Поля коннекторов — Коннекторы; выбор DataFlow vs DataFlowCron — Типы нагрузки.
См. также: FAQ · Best Practices · DataFlowCron · Helm Values · CLI
Пути
- Установка оператора — Helm / CRD
- Первый потоковый конвейер —
DataFlowсSecretRef - Первый cron-конвейер —
DataFlowCron - Локальная разработка — Go, Docker, kind
Установка оператора
Предварительные требования
- Kubernetes-кластер (1.24+)
- Helm 3.0+
- kubectl, настроенный на кластер
- Доступ к источникам данных (Kafka, PostgreSQL и т.д.)
Управление CRD
Оператор ставит два CRD: DataFlow (dataflows) и DataFlowCron (dataflowcrons).
Автоматическая установка (через Helm)
При установке через Helm (рекомендуется) CRD ставятся и обновляются автоматически при crds.install: true (по умолчанию). Отдельный kubectl apply не нужен.
Ручная установка
Если CRD управляются отдельно (ArgoCD, FluxCD или crds.install: false), примените оба CRD:
kubectl apply -f https://raw.githubusercontent.com/dataflow-operator/dataflow/refs/heads/main/config/crd/bases/dataflow.dataflow.io_dataflows.yaml
kubectl apply -f https://raw.githubusercontent.com/dataflow-operator/dataflow/refs/heads/main/config/crd/bases/dataflow.dataflow.io_dataflowcrons.yaml
Или из локального checkout:
kubectl apply -f dataflow/config/crd/bases/dataflow.dataflow.io_dataflows.yaml
kubectl apply -f dataflow/config/crd/bases/dataflow.dataflow.io_dataflowcrons.yaml
Конфигурация CRD в Helm
| Параметр | По умолчанию | Описание |
|---|---|---|
crds.install |
true |
Устанавливать и обновлять CRD при helm install / helm upgrade |
crds.keep |
true |
Аннотация helm.sh/resource-policy: keep — CRD не удаляются при helm uninstall |
Обновление: CRD обновляются при каждом helm upgrade.
Удаление: При crds.keep: true (по умолчанию) CRD остаются после helm uninstall. Чтобы отключить установку CRD через Helm:
crds:
install: false
Установка через Helm (рекомендуется)
Базовая установка
helm install dataflow-operator oci://ghcr.io/dataflow-operator/helm-charts/dataflow-operator
Установка в namespace default с настройками чарта по умолчанию.
Конкретный namespace
helm install dataflow-operator oci://ghcr.io/dataflow-operator/helm-charts/dataflow-operator \
--namespace dataflow-system \
--create-namespace
Для локальной разработки чарта:
helm install dataflow-operator ./helm-charts/dataflow-operator
Кастомные настройки
helm install dataflow-operator oci://ghcr.io/dataflow-operator/helm-charts/dataflow-operator \
--set image.repository=your-registry/controller \
--set image.tag=v1.0.0 \
--set replicaCount=2 \
--set resources.limits.memory=1Gi \
--set resources.limits.cpu=500m \
--set resources.requests.memory=256Mi \
--set resources.requests.cpu=100m
Файл values
image:
repository: your-registry/controller
tag: v1.0.0
replicaCount: 2
resources:
limits:
cpu: 1000m
memory: 1Gi
requests:
cpu: 500m
memory: 512Mi
serviceAccount:
create: true
annotations:
eks.amazonaws.com/role-arn: arn:aws:iam::123456789012:role/dataflow-operator
securityContext:
runAsNonRoot: true
runAsUser: 1000
fsGroup: 1000
# Опционально: Sentry
# sentry:
# enabled: true
# dsn: "https://xxx@o0.ingest.sentry.io/123"
# environment: production
# tracesSampleRate: 0.1
helm install dataflow-operator oci://ghcr.io/dataflow-operator/helm-charts/dataflow-operator -f my-values.yaml
Проверка
kubectl get pods -l app.kubernetes.io/name=dataflow-operator
# Оба CRD
kubectl get crd dataflows.dataflow.dataflow.io dataflowcrons.dataflow.dataflow.io
kubectl logs -l app.kubernetes.io/name=dataflow-operator --tail=50
kubectl get deployment dataflow-operator
Ожидаемый вывод:
NAME READY STATUS RESTARTS AGE
dataflow-operator-7d8f9c4b5d-xxxxx 1/1 Running 0 1m
Обновление
helm upgrade dataflow-operator oci://ghcr.io/dataflow-operator/helm-charts/dataflow-operator
С values:
helm upgrade dataflow-operator oci://ghcr.io/dataflow-operator/helm-charts/dataflow-operator -f my-values.yaml
Конкретная версия:
helm upgrade dataflow-operator oci://ghcr.io/dataflow-operator/helm-charts/dataflow-operator \
--set image.tag=v1.1.0
Удаление
helm uninstall dataflow-operator
При crds.keep: true CRD остаются; существующие ресурсы перестают реконсилироваться.
Чтобы удалить CRD и все ресурсы DataFlow / DataFlowCron:
kubectl delete dataflow --all --all-namespaces
kubectl delete dataflowcron --all --all-namespaces
helm uninstall dataflow-operator
kubectl delete crd dataflows.dataflow.dataflow.io
kubectl delete crd dataflowcrons.dataflow.dataflow.io
Warning
Удаление CRD удаляет все ресурсы DataFlow и DataFlowCron во всех namespace.
Первый потоковый конвейер
Непрерывный DataFlow (Deployment). В кластере / production для credentials используйте SecretRef.
Kafka → PostgreSQL с SecretRef
apiVersion: v1
kind: Secret
metadata:
name: kafka-credentials
namespace: default
type: Opaque
stringData:
brokers: "kafka-broker:9092"
topic: "input-topic"
consumerGroup: "dataflow-group"
---
apiVersion: v1
kind: Secret
metadata:
name: postgres-credentials
namespace: default
type: Opaque
stringData:
connectionString: "postgres://user:password@postgres-host:5432/dbname?sslmode=require"
table: "output_table"
---
apiVersion: dataflow.dataflow.io/v1
kind: DataFlow
metadata:
name: kafka-to-postgres
namespace: default
spec:
source:
type: kafka
config:
brokersSecretRef:
name: kafka-credentials
key: brokers
topicSecretRef:
name: kafka-credentials
key: topic
consumerGroupSecretRef:
name: kafka-credentials
key: consumerGroup
sink:
type: postgresql
config:
connectionStringSecretRef:
name: postgres-credentials
key: connectionString
tableSecretRef:
name: postgres-credentials
key: table
autoCreateTable: true
Sample из репозитория (включая SASL SecretRef):
kubectl apply -f dataflow/config/samples/kafka-to-postgres-secrets.yaml
Подробнее: Использование Secrets. Ресурсы и placement подов: Примеры.
Проверка статуса
kubectl get dataflow kafka-to-postgres
kubectl describe dataflow kafka-to-postgres
kubectl get dataflow kafka-to-postgres -o yaml
Пример status:
status:
phase: Running
processedCount: 150
errorCount: 0
lastProcessedTime: "2024-01-15T10:30:00Z"
message: "Processing messages successfully"
Отправка тестового сообщения
kafka-console-producer --broker-list kafka-broker:9092 --topic input-topic
# Введите JSON и нажмите Enter
{"id": 1, "name": "Test", "value": 100}
Или (локальный стек):
./scripts/send-test-message.sh
Проверка PostgreSQL
psql "$CONNECTION_STRING" -c "SELECT * FROM output_table;"
Первый cron-конвейер
DataFlowCron запускает тот же Source → Transform → Sink по cron-расписанию (CronJob / Job), с опциональными post-run триггерами.
Лучше всего подходит для polling / batch источников (postgresql, trino, clickhouse, nessie, iceberg). Kafka — streaming и обычно не завершается по исчерпанию источника — см. DataFlowCron и Типы нагрузки.
Минимальный пример (PostgreSQL → PostgreSQL)
apiVersion: v1
kind: Secret
metadata:
name: postgres-credentials
namespace: default
type: Opaque
stringData:
connectionString: "postgres://user:password@postgres:5432/db?sslmode=require"
---
apiVersion: dataflow.dataflow.io/v1
kind: DataFlowCron
metadata:
name: pg-nightly-sync
spec:
schedule: "0 2 * * *"
concurrencyPolicy: Forbid
source:
type: postgresql
config:
connectionStringSecretRef:
name: postgres-credentials
key: connectionString
table: public.orders_staging
pollInterval: 5
sink:
type: postgresql
config:
connectionStringSecretRef:
name: postgres-credentials
key: connectionString
table: public.orders_warehouse
autoCreateTable: true
batchSize: 200
kubectl apply -f pg-nightly-sync.yaml
kubectl get dataflowcron pg-nightly-sync
kubectl get cronjob,job -l dataflow.dataflow.io/dataflow-cron=pg-nightly-sync
Sample в репозитории (Kafka → Nessie + triggers): dataflow/config/samples/dataflowcron-example.yaml. Больше YAML: Примеры DataFlowCron.
Локальная разработка
Для контрибьюторов: оператор против docker-compose / kind. Plaintext connection strings допустимы только в этом локальном пути.
Предварительные требования
- Go 1.25+
- Docker и Docker Compose
- Task (опционально)
- Порты: 8080, 5050, 15672, 8081, 5432, 9092, 5672
Запуск зависимостей
docker-compose up -d
Запускает Kafka (9092) + Kafka UI (8080), PostgreSQL (5432) + pgAdmin (5050).
- Kafka UI: http://localhost:8080
- pgAdmin: http://localhost:5050 —
admin@admin.com/admin
DataFlow только для local (plaintext)
Только локальная разработка
Не используйте plaintext credentials в shared или production-кластерах. Предпочтительнее SecretRef.
apiVersion: dataflow.dataflow.io/v1
kind: DataFlow
metadata:
name: kafka-to-postgres
namespace: default
spec:
source:
type: kafka
config:
brokers:
- localhost:9092
topic: input-topic
consumerGroup: dataflow-group
sink:
type: postgresql
config:
connectionString: "postgres://dataflow:dataflow@localhost:5432/dataflow?sslmode=disable"
table: output_table
autoCreateTable: true
kubectl apply -f dataflow/config/samples/kafka-to-postgres.yaml
Запуск оператора локально
task install # CRD (kind / minikube)
task run
Или:
./scripts/run-local.sh
Опционально: kind-кластер
./scripts/setup-kind.sh
task install
task run
Отладка
kubectl logs -l app.kubernetes.io/name=dataflow-operator -f
kubectl get events --sort-by='.lastTimestamp' | grep dataflow
Следующие шаги
- Коннекторы — источники и приёмники
- Трансформации — преобразования сообщений
- Примеры — практические манифесты
- Agent Skills — AI-помощь с deploy и манифестами
- Разработка — участие в разработке
Устранение неполадок
Оператор не запускается
kubectl logs -l app.kubernetes.io/name=dataflow-operator
kubectl describe pod -l app.kubernetes.io/name=dataflow-operator
kubectl get crd dataflows.dataflow.dataflow.io dataflowcrons.dataflow.dataflow.io -o yaml
DataFlow не обрабатывает сообщения
kubectl describe dataflow <name>- Проверьте доступность источника (Kafka consumer / запрос
psql) - Логи оператора / процессора
Проблемы с подключением
- Доступность из кластера (NetworkPolicies, DNS)
- Credentials и ключи Secret
- Локально:
localhostилиhost.docker.internal
Trino: REMOTE_TASK_ERROR (код 65542)
Воркер не ответил координатору. Оператор ретраит (до 5 попыток с backoff).
Если ошибки продолжаются: уменьшите batchSize Trino sink (например, 3–5), проверьте воркеры Trino, увеличьте таймауты запросов.