Skip to content

Commit 172cf0f

Browse files
committed
Refactoring
1 parent 66d0877 commit 172cf0f

43 files changed

Lines changed: 1752 additions & 2350 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎README.md‎

Lines changed: 93 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -1,43 +1,79 @@
11
# DataFlow Operator
22

3-
Kubernetes operator for streaming data between different sources (Kafka, PostgreSQL, Trino) with support for message transformations.
3+
Kubernetes operator for streaming data between different sources and sinks with support for message transformations.
44

55
**[Online documentation](https://dataflow-operator.github.io/docs/)**
66

7-
## Quick Start
7+
## Architecture
8+
9+
The operator watches `DataFlow` custom resources and creates processor pods. Each processor reads from a source, applies optional transformations, and writes to a sink. Configuration uses the `type` + `config` format:
10+
11+
```yaml
12+
apiVersion: dataflow.dataflow.io/v1
13+
kind: DataFlow
14+
metadata:
15+
name: kafka-to-postgres
16+
spec:
17+
source:
18+
type: kafka
19+
config:
20+
brokers:
21+
- localhost:9092
22+
topic: input-topic
23+
consumerGroup: dataflow-group
24+
sink:
25+
type: postgresql
26+
config:
27+
connectionString: "postgres://user:pass@localhost:5432/db?sslmode=disable"
28+
table: output_table
29+
```
30+
31+
### Supported connectors
32+
33+
| Source | Sink |
34+
|--------|------|
35+
| Kafka | Kafka |
36+
| PostgreSQL | PostgreSQL |
37+
| — | ClickHouse |
38+
| — | Trino |
39+
| — | Nessie (Apache Iceberg) |
840
9-
### Prerequisites
41+
### Supported transformations
42+
43+
Filter, Select, Remove, Mask, Flatten, Timestamp, SnakeCase, Router, Chain.
44+
45+
## Prerequisites
1046
1147
- Kubernetes 1.24+
1248
- Helm 3.0+
1349
- kubectl
14-
- Go 1.21+ (for local development)
50+
- Go 1.25+ (for local development)
1551
- Docker and docker-compose (for local development)
1652
17-
#### Local Development
53+
## Quick Start
1854
19-
For local development, you can run the operator locally:
20-
```bash
21-
task run
22-
```
55+
### Install via Helm
2356
24-
Or use the script:
2557
```bash
26-
./scripts/run-local.sh
58+
helm repo add dataflow-operator https://dataflow-operator.github.io/helm-charts
59+
helm repo update
60+
helm install dataflow-operator dataflow-operator/dataflow-operator
2761
```
2862

2963
### Local Development Setup
3064

31-
1. Start dependencies with UI interfaces:
65+
1. Start dependencies:
66+
3267
```bash
3368
docker-compose up -d
3469
```
3570

3671
Available UIs:
3772
- **Kafka UI**: http://localhost:8080
38-
- **pgAdmin**: http://localhost:5050 (admin@admin.com / admin)
73+
- **ClickHouse**: http://localhost:8123
3974

4075
2. Run the operator:
76+
4177
```bash
4278
task run
4379
```
@@ -46,27 +82,65 @@ task run
4682

4783
### Code Generation
4884

49-
If you encounter issues with `task generate`, try:
85+
```bash
86+
task generate # DeepCopy, CRD, RBAC manifests
87+
task manifests # CRD and RBAC only
88+
```
89+
90+
If you encounter issues with `task generate`:
5091

5192
```bash
52-
# Update controller-gen
5393
go install sigs.k8s.io/controller-tools/cmd/controller-gen@latest
54-
55-
# Then
5694
task generate
5795
```
5896

97+
### Building
98+
99+
```bash
100+
task build # builds bin/manager
101+
```
102+
103+
Docker image (builds both operator and processor binaries):
104+
105+
```bash
106+
docker build -t dataflow:local .
107+
```
108+
59109
### Testing
60110

61111
```bash
62112
# Unit tests
63113
task test
64114

65-
# Integration tests (requires kind)
66-
./scripts/setup-kind.sh
115+
# Integration tests (requires Docker — uses testcontainers)
67116
task test-integration
68117
```
69118

119+
### Taskfile commands
120+
121+
| Command | Description |
122+
|---------|-------------|
123+
| `task build` | Build `bin/manager` |
124+
| `task run` | Run operator locally |
125+
| `task generate` | Generate DeepCopy and manifests |
126+
| `task manifests` | Generate CRD and RBAC |
127+
| `task test` | Unit tests |
128+
| `task test-integration` | Integration tests (Docker required) |
129+
| `task install` | Install CRDs into cluster |
130+
| `task uninstall` | Remove CRDs from cluster |
131+
132+
## Configuration samples
133+
134+
Example manifests are in `config/samples/`:
135+
136+
- `kafka-to-postgres.yaml` — Kafka to PostgreSQL
137+
- `kafka-to-clickhouse.yaml` — Kafka to ClickHouse
138+
- `kafka-to-trino.yaml` — Kafka to Trino
139+
- `postgres-to-kafka-router.yaml` — PostgreSQL to Kafka with Router
140+
- `clickhouse-to-clickhouse.yaml` — ClickHouse to ClickHouse
141+
142+
See all samples: `ls config/samples/`
143+
70144
## License
71145

72146
Apache License 2.0

‎Taskfile.yml‎

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,50 @@ tasks:
9696
deps: [manifests, generate, fmt, vet]
9797
cmds:
9898
- go run ./main.go
99+
sync-crd-to-helm:
100+
desc: Sync generated CRD into Helm chart template (wraps CRD with Helm conditionals)
101+
deps: [manifests]
102+
vars:
103+
CRD_SOURCE: config/crd/bases/dataflow.dataflow.io_dataflows.yaml
104+
HELM_TARGET: ../helm-charts/charts/dataflow-operator/templates/crd.yaml
105+
cmds:
106+
- |
107+
HEADER='{{- if .Values.crds.install }}'
108+
FOOTER='{{- end }}'
109+
LABELS_BLOCK=' labels:
110+
{{- include "dataflow-operator.labels" . | nindent 4 }}'
111+
KEEP_BLOCK=' {{- if .Values.crds.keep }}
112+
"helm.sh/resource-policy": keep
113+
{{- end }}'
114+
115+
{
116+
echo "$HEADER"
117+
awk '
118+
/^---$/ { next }
119+
/^metadata:$/ {
120+
print
121+
getline
122+
# annotations block
123+
if ($0 ~ /^ annotations:/) {
124+
print
125+
# print existing annotations
126+
while (getline > 0 && $0 ~ /^ [^ ]/) {
127+
print
128+
}
129+
# insert keep annotation block
130+
printf "%s\n", "'"$KEEP_BLOCK"'"
131+
# insert labels block
132+
printf "%s\n", "'"$LABELS_BLOCK"'"
133+
print
134+
}
135+
next
136+
}
137+
{ print }
138+
' "{{.CRD_SOURCE}}"
139+
echo "$FOOTER"
140+
} > "{{.HELM_TARGET}}"
141+
echo "Synced CRD to {{.HELM_TARGET}}"
142+
99143
# Deployment tasks
100144
install:
101145
desc: Install CRDs into the K8s cluster specified in ~/.kube/config
@@ -147,6 +191,11 @@ tasks:
147191
preconditions:
148192
- test -d {{.LOCALBIN}} || mkdir -p {{.LOCALBIN}}
149193

194+
test-helm-crd:
195+
desc: Run Helm CRD template tests (requires helm CLI)
196+
cmds:
197+
- bash ../helm-charts/charts/dataflow-operator/tests/crd_test.sh
198+
150199
tidy:
151200
desc: Run go mod tidy
152201
cmds:

0 commit comments

Comments
 (0)