Kafka Deployment (MODE 3)
Status: Preview
Overview
Use MODE 3 when:
-
Kafka infrastructure already exists in your environment
-
Security requirements demand events not be written to disk (credit cards, PII, etc.)
-
You need encrypted transport (SSL/SASL_SSL)
-
Direct stream processing is preferred over log-based ingestion
-
You want to leverage Kafka’s at-least-once delivery guarantees
Event Pipeline
Quarkus Flow
→ Kafka (CloudEvents, topic: flow-lifecycle-out)
→ Data Index Ingestion Service
→ PostgreSQL (workflow_instances, task_instances)
→ Data Index GraphQL API
(failed records → data-index-events-dlq)
Kafka Topic Configuration
Required Topics
-
flow-lifecycle-out- Main event topic (published by Quarkus Flow applications) -
data-index-events-dlq- Dead-letter queue for failed records
|
You can change topic names via configuration through environment variables:
|
Creating Topics
In non-production environments, topics are typically auto-created. In production, create them explicitly in accordance with your Kafka cluster management practices.
# Main topic (replicas=3, partitions=3)
kafka-topics.sh --create \
--bootstrap-server kafka.kafka.svc.cluster.local:9092 \
--topic flow-lifecycle-out \
--replication-factor 3 \
--partitions 3 \
--config retention.ms=86400000 \
--config min.insync.replicas=2
# DLQ topic
kafka-topics.sh --create \
--bootstrap-server kafka.kafka.svc.cluster.local:9092 \
--topic data-index-events-dlq \
--replication-factor 3 \
--partitions 1 \
--config retention.ms=604800000
Kubernetes Deployment
Basic Manifest
apiVersion: apps/v1
kind: Deployment
metadata:
name: data-index-ingestion-kafka
namespace: data-index
spec:
replicas: 1
selector:
matchLabels:
app: data-index-ingestion-kafka
template:
metadata:
labels:
app: data-index-ingestion-kafka
spec:
containers:
- name: kafka-ingestion
image: kubesmarts/data-index-ingestion-kafka-service:999-SNAPSHOT
ports:
- containerPort: 8080
name: http
env:
- name: KAFKA_BOOTSTRAP_SERVERS
value: "kafka.kafka.svc.cluster.local:9092"
- name: QUARKUS_DATASOURCE_JDBC_URL
value: "jdbc:postgresql://postgresql:5432/data-index"
- name: QUARKUS_DATASOURCE_USERNAME
valueFrom:
secretKeyRef:
name: database-credentials
key: username
- name: QUARKUS_DATASOURCE_PASSWORD
valueFrom:
secretKeyRef:
name: database-credentials
key: password
livenessProbe:
httpGet:
path: /q/health/live
port: 8080
initialDelaySeconds: 30
periodSeconds: 10
readinessProbe:
httpGet:
path: /q/health/ready
port: 8080
initialDelaySeconds: 10
periodSeconds: 5
resources:
requests:
cpu: 500m
memory: 512Mi
limits:
cpu: 2000m
memory: 2Gi
---
apiVersion: v1
kind: Service
metadata:
name: data-index-ingestion-kafka
namespace: data-index
spec:
selector:
app: data-index-ingestion-kafka
ports:
- port: 8080
targetPort: 8080
name: http
Configuration
Required Environment Variables
| Variable | Default | Description |
|---|---|---|
|
localhost:29092 |
Kafka broker URLs (comma-separated) |
|
jdbc:h2:mem:test |
PostgreSQL JDBC connection string |
|
(dev services) |
Database username |
|
(dev services) |
Database password |
Optional Configuration
| Variable | Default | Description |
|---|---|---|
|
|
Kafka topic name |
|
|
Consumer group |
|
|
Enable channel health checks |
|
|
Include channel in readiness checks |
|
|
DLQ topic name |
Monitoring
Event Processing
Field-Level Idempotency
MODE 3 guarantees idempotency for out-of-order and duplicate events:
- Immutable fields (first value wins)
-
-
start, input, name, version, namespace
-
Never updated after initial insertion
-
- Terminal fields (last non-null wins)
-
-
end, output, error fields
-
Updated only if incoming event timestamp is newer
-
- Status precedence
-
-
COMPLETED, FAULTED, CANCELLED > RUNNING > CREATED
-
Terminal states override less-terminal states
-
Out-of-Order Recovery
If a task event arrives before the parent workflow:
-
Task event consumed → INSERT fails (foreign key constraint)
-
Savepoint rolled back
-
Placeholder workflow created with minimal data
-
Task event retried → INSERT succeeds
-
Workflow event arrives later → updates placeholder with full data
This ensures no task events are lost due to event ordering.
Troubleshooting
Service won’t start
Check logs:
kubectl logs deployment/data-index-ingestion-kafka -n data-index
Common causes:
* PostgreSQL unreachable → verify QUARKUS_DATASOURCE_JDBC_URL
* Kafka unreachable → verify KAFKA_BOOTSTRAP_SERVERS
* Database schema missing → run Flyway migrations
Events not consumed
Check readiness:
kubectl get pods -n data-index | grep data-index-ingestion-kafka
# Check logs
kubectl logs deployment/data-index-ingestion-kafka -n data-index | grep -i error
DLQ messages pile up
Inspect failed events:
kafka-console-consumer.sh \
--bootstrap-server kafka:9092 \
--topic data-index-events-dlq \
--max-messages 5 | jq
Common causes: * Malformed CloudEvents → fix event publisher * Database unavailable → events will retry once recovered * Schema mismatch → upgrade service or downgrade event publisher
Comparison
| Feature | MODE 1 (Vector + Triggers) | MODE 2 (Vector + ES) | MODE 3 (Kafka) |
|---|---|---|---|
Event Source |
Container logs |
Container logs |
Kafka topics |
Ingestion |
Vector DaemonSet |
Vector DaemonSet |
SmallRye Reactive Messaging |
Normalization |
PostgreSQL triggers |
ES transforms |
Java processors (JDBC) |
Raw Storage |
|
|
None (direct to normalized) |
Performance |
~10ms latency |
~1s latency |
~100ms latency |
DLQ |
N/A |
N/A |
Yes ( |
Security |
Disk files |
Disk files |
Kafka (SSL/SASL capable) |
End-to-End Testing
Local KIND Cluster Test
Run the complete MODE 3 e2e test suite to verify the entire Kafka ingestion pipeline:
cd data-index/scripts/kind
./test-mode3-e2e.sh
What this tests:
-
✓ KIND cluster creation and configuration
-
✓ PostgreSQL deployment (normalized storage)
-
✓ Kafka broker deployment (KRaft single-node)
-
✓ Data Index query service deployment (PostgreSQL backend)
-
✓ Kafka ingestion service deployment
-
✓ Workflow test app deployment (Kafka profile)
-
✓ Test workflow execution via REST API
-
✓ CloudEvents published to Kafka topic
-
✓ Ingestion service consumption from Kafka
-
✓ Event normalization to PostgreSQL
-
✓ GraphQL API returning normalized data
-
✓ Idempotency (replaying events doesn’t create duplicates)
Pipeline flow verified:
workflow-test-app (Kafka profile)
↓
Kafka topic: flow-lifecycle-out (CloudEvents)
↓
data-index-ingestion-kafka-service
↓
PostgreSQL (workflow_instances, task_instances)
↓
Data Index GraphQL API
Test output example:
[INFO] ==========================================
[INFO] MODE 3 (Kafka) End-to-End Integration Test
[INFO] ==========================================
[STEP] Creating KIND cluster...
[INFO] ✓ Cluster ready
[STEP] Creating namespaces...
[INFO] ✓ Namespaces ready
[STEP] Installing PostgreSQL...
[INFO] ✓ PostgreSQL ready
[STEP] Schema will be applied by Flyway on ingestion service startup
[INFO] ✓ Schema initialization delegated to Flyway
[STEP] Installing Kafka (KRaft single-node)...
[INFO] Waiting for Kafka to be ready (this may take ~60 seconds)...
[INFO] ✓ Kafka broker accepting connections
[INFO] ✓ Kafka ready at kafka.kafka.svc.cluster.local:9092
[STEP] Deploying data-index-service (postgresql query backend)...
[INFO] ✓ Query service ready at http://localhost:30080/graphql
[STEP] Deploying data-index-ingestion-kafka-service...
[INFO] ✓ Ingestion service ready and consuming from Kafka
[STEP] Deploying workflow-test-app (Kafka profile)...
[INFO] ✓ workflow-test-app ready (publishing to Kafka)
[STEP] Executing test workflows via REST API...
[INFO] Triggering simple-set workflow...
[INFO] → HTTP 200: simple-set workflow triggered
[INFO] ✓ Workflows triggered — events are now in Kafka topic 'flow-lifecycle-out'
[STEP] Verifying events in Kafka topic 'flow-lifecycle-out'...
[INFO] ✓ Found 10 messages in flow-lifecycle-out
[STEP] Verifying normalized data in PostgreSQL...
[INFO] ✓ Found 2 normalized workflow instance(s)
[INFO] Task instances: 4
[INFO] Sample workflow instance:
id | name | status | has_start
--------------------------------------+-------------+-----------+-----------
01KQ7X... | simple-set | COMPLETED | t
[INFO] ✓ Normalization verified
[STEP] Verifying GraphQL API returns normalized data...
[INFO] Sample workflow from GraphQL: 01KQ7X...
[INFO] ✓ GraphQL API verified
[STEP] Verifying idempotency (re-triggering same workflow)...
[INFO] Workflow count before re-trigger: 2
[INFO] Workflow count after re-trigger: 2
[INFO] ✓ Idempotency verified (no unexpected duplicates)
[INFO] ==========================================
[INFO] MODE 3 (Kafka) E2E Test Complete!
[INFO] ==========================================
Pipeline:
workflow-test-app → Kafka (flow-lifecycle-out) → ingestion-service → PostgreSQL → GraphQL
Results:
✓ Kafka broker running
✓ Workflow events published as CloudEvents
✓ Ingestion service consumed and normalized events
✓ PostgreSQL: 2 workflow instance(s), 4 task instance(s)
✓ GraphQL API returning data
Access Points:
GraphQL API: http://localhost:30080/graphql
GraphQL UI: http://localhost:30080/q/graphql-ui
PostgreSQL: postgresql://dataindex:dataindex123@localhost:30432/dataindex
Kafka: localhost:30900 (NodePort, for tools like kcat)
[INFO] ✅ All MODE 3 tests passed!
Manual verification steps:
# Watch live Kafka events
kubectl exec -n kafka kafka-0 -- \
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic flow-lifecycle-out \
--from-beginning
# Check consumer group lag
kubectl exec -n kafka kafka-0 -- \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server localhost:9092 \
--describe --group data-index-ingestion
# Query PostgreSQL directly
kubectl exec -n postgresql postgresql-0 -- \
env PGPASSWORD=dataindex123 psql -U dataindex -d dataindex \
-c "SELECT id, name, status FROM workflow_instances;"
# Query GraphQL API
curl http://localhost:30080/graphql \
-H "Content-Type: application/json" \
-d '{"query":"{ getWorkflowInstances { id name status } }"}'
Test duration: ~6-8 minutes (includes cluster creation, Kafka startup, event publishing and consumption)
CloudEvents verification:
The test verifies that workflow events are published as CloudEvents to Kafka:
# Inspect a CloudEvent
kubectl exec -n kafka kafka-0 -- \
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic flow-lifecycle-out \
--max-messages 1 | jq
# Expected structure:
{
"specversion": "1.0",
"type": "io.serverlessworkflow.workflow.started.v1",
"source": "workflow-test-app",
"id": "01KQ7X...",
"time": "2026-08-31T10:30:00Z",
"datacontenttype": "application/json",
"data": {
"workflowInstanceId": "01KQ7X...",
"workflowName": "simple-set",
"status": "RUNNING",
...
}
}
Cleanup:
# Delete the test cluster
kind delete cluster --name data-index-test
See Local Development with KIND for manual setup steps and troubleshooting.