Skip to content

Repository files navigation

DataFlow Operator

A production-ready Kubernetes operator for streaming data pipelines with support for multiple sources, sinks, and comprehensive message transformations.

📖 Full Documentation


Overview

DataFlow Operator automates the deployment and management of data streaming pipelines in Kubernetes. It watches custom DataFlow resources and orchestrates processor pods that read from sources, apply optional transformations, and write to sinks. The operator handles fault tolerance, checkpointing, scheduling, and comprehensive monitoring out of the box.

Key Features

  • Multi-Source/Sink Support: Kafka, PostgreSQL, ClickHouse, Trino, and Apache Iceberg (Nessie). Securely configure connectors using Kubernetes Secrets via SecretRef fields.
  • Rich Transformations: Filter, Select, Remove, Mask, Flatten, Timestamp, SnakeCase, Router, Chain.
  • Web GUI: Built-in browser-based interface to manage manifests, stream logs, and monitor metrics without kubectl.
  • MCP Server: Model Context Protocol server for AI agents to validate, generate, and migrate DataFlow manifests.
  • Agent Skills: Portable AI guides (github.com/dataflow-operator/skills, docs).
  • Error Handling: Route failed messages to a dedicated error sink (e.g., Dead Letter Queue) with full error context and original payloads to prevent data loss.
  • Fault Tolerance: At-least-once delivery semantics with checkpoint persistence, graceful shutdown, and support for idempotent sinks (e.g., UPSERT, ReplacingMergeTree) to safely handle duplicates.
  • Scheduled Pipelines: DataFlowCron for time-based pipeline execution with triggers.
  • High Availability: Leader election for multi-replica deployments.
  • Observable: Prometheus metrics, structured logging, Sentry integration, health probes, and native Kubernetes Events for lifecycle audit trails.
  • Kubernetes Native: Custom Resource Definitions, RBAC, Helm charts.

Architecture

How It Works

  1. Operator Controller: Watches DataFlow and DataFlowCron resources in your cluster.
  2. Processor Pod: Creates ephemeral or long-running pods that execute the data pipeline.
  3. Pipeline Execution: source → transformations → sink flow with built-in error handling and error sink routing.
  4. State Management: Stores checkpoint data in ConfigMaps (or native offsets for Kafka) for recovery.
  5. Observability: Emits standard Kubernetes Events (Normal/Warning) for reconciliation lifecycle auditing.

Configuration Example

apiVersion: dataflow.dataflow.io/v1kind: DataFlowmetadata:
name: kafka-to-postgresspec:
source:
type: kafkaconfig:
brokers:
- kafka-broker:9092 topic: input-eventsconsumerGroup: dataflow-grouptransformations:
- type: filterconfig:
expression: "event_type == 'purchase'"
- type: maskconfig:
fields:
- credit_cardsink:
type: postgresqlconfig:
connectionStringSecretRef: name: db-credentialskey: urltable: events# Error sink catches failed writes (e.g. constraint violations) and saves them for replayerrors:
type: postgresqlconfig:
connectionStringSecretRef: name: db-credentialskey: urltable: events_dead_letterautoCreateTable: true

Quick Start

1. Install via Helm

# Add the Helm repository
helm repo add dataflow-operator https://dataflow-operator.github.io/helm-charts
helm repo update
# Install the operator
helm install dataflow-operator dataflow-operator/dataflow-operator \
--namespace dataflow-system \
--create-namespace

2. Deploy Your First Pipeline

# Create a sample Kafka-to-PostgreSQL pipeline
kubectl apply -f - <<EOFapiVersion: dataflow.dataflow.io/v1kind: DataFlowmetadata: name: my-pipeline namespace: defaultspec: source: type: kafka config: brokers: - kafka:9092 topic: events consumerGroup: my-app sink: type: postgresql config: connectionString: "postgres://user:pass@postgres:5432/db" table: eventsEOF# Monitor the pipeline status and Kubernetes events
kubectl describe dataflow my-pipeline
kubectl get events --watch

3. Local Development Setup

# Start local infrastructure (Kafka, PostgreSQL, ClickHouse)
docker-compose up -d
# Available UIs:# - Kafka UI: http://localhost:8080# - ClickHouse: http://localhost:8123# Run the operator locally
task run
# In another terminal, apply a sample configuration
kubectl apply -f config/samples/kafka-to-postgres.yaml

Resources

About

DataFlow Operator is a Kubernetes operator for streaming data between different data sources with support for message transformations.

Topics

Resources

Stars

3 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages