Kafka Compacted Topic Explorer
KGazer is a developer tool for exploring and inspecting Kafka compacted topics. It continuously consumes messages from your Kafka clusters, stores them in PostgreSQL and provides a web interface to browse keys, view message history and compare changes over time.
If you've ever needed to answer "what's the current value for this key?" or "what changed in this key's history?", KGazer gives you that visibility without writing throwaway consumer scripts.
Features
- Multi-cluster support — Connect to multiple Kafka clusters simultaneously
- Key browser — Search, filter and paginate through all keys in a compacted topic
- Key value search — Search message values using
field:valuesyntax with field autocomplete - Message history — View every version of a key with syntax-highlighted JSON and inline diffs
- Change timeline — Interactive visual timeline showing when a key changed and how much
- Consumer group monitoring — See which consumer groups are reading a topic, per-partition lag and consumer assignments
- Offset management — Reset consumer group offsets to earliest, latest or specific per-partition values
- Avro & Protobuf support — Automatic deserialisation via Schema Registry (Confluent-compatible)
- Topic lifecycle — Detects new topics, cleans up deleted topics and handles topic recreation transparently
- Re-consume — Wipe stored data and re-consume a topic from the beginning with one click
- Kubernetes deployment — Helm chart with bundled PostgreSQL or external database support
Quick Start
Prerequisites
- Docker and Docker Compose
1. Configure your clusters
cp config.yml.example config.yml
Edit config.yml with your Kafka cluster details:
kafka: clusters: - name: my-cluster bootstrapServers: localhost:9092 - name: production bootstrapServers: kafka.example.com:9092 properties: security.protocol: SASL_SSL sasl.mechanism: PLAIN sasl.jaas.config: >- org.apache.kafka.common.security.plain.PlainLoginModule required username="user" password="pass"; schemaRegistry: https://schema-registry.example.com schemaRegistryAuth: username: sr-user password: sr-pass kgazer: server: port: 8080 db: host: localhost port: 5432 name: kgazer user: kgazer password: kgazer sslmode: disable compactedOnly: true
Cluster names must be URL-safe: letters, digits, hyphens, dots and underscores only.
2. Start the stack
This starts three services:
| Service | Default Port | Description |
|---|---|---|
| frontend | localhost:4200 | React web interface with hot reload |
| backend | localhost:1899 | Go API server |
| db | localhost:5899 | PostgreSQL 16 |
Open http://localhost:4200 and KGazer will begin discovering and consuming topics automatically.
Ports are configurable via environment variables:
KGAZER_PORT_UI=3000 KGAZER_PORT_API=9090 KGAZER_PORT_DB=5433 docker compose up
3. Run the tests
docker compose --profile test run --rm test # Backend (Go) docker compose --profile test run --rm test-frontend # Frontend (Vitest)
Deploying to Kubernetes
KGazer ships with a Helm chart in charts/kgazer/. The chart deploys the backend, a bundled PostgreSQL instance (via the Bitnami subchart) and an optional Ingress.
Prerequisites
- A Kubernetes cluster (1.25+)
- Helm 3.x
Install
helm dependency update charts/kgazer helm install kgazer charts/kgazer \ --set 'kafka.clusters[0].name=my-cluster' \ --set 'kafka.clusters[0].bootstrapServers=kafka.example.com:9092'
This deploys KGazer with a bundled PostgreSQL database. To verify the installation:
kubectl port-forward svc/kgazer 8080:8080 curl http://localhost:8080/api/health
Using an existing PostgreSQL database
To connect to an existing PostgreSQL instance instead of deploying one:
helm install kgazer charts/kgazer \ --set postgresql.enabled=false \ --set externalDatabase.enabled=true \ --set externalDatabase.host=my-postgres.example.com \ --set externalDatabase.port=5432 \ --set externalDatabase.name=kgazer \ --set externalDatabase.user=kgazer \ --set externalDatabase.password=changeme \ --set externalDatabase.sslmode=require \ --set 'kafka.clusters[0].name=my-cluster' \ --set 'kafka.clusters[0].bootstrapServers=kafka.example.com:9092'
If the database password is already stored in a Kubernetes Secret, reference it directly:
helm install kgazer charts/kgazer \ --set postgresql.enabled=false \ --set externalDatabase.enabled=true \ --set externalDatabase.host=my-postgres.example.com \ --set externalDatabase.name=kgazer \ --set externalDatabase.user=kgazer \ --set externalDatabase.existingSecret=my-db-credentials \ --set 'kafka.clusters[0].name=my-cluster' \ --set 'kafka.clusters[0].bootstrapServers=kafka.example.com:9092'
The Secret must contain a password key.
Configuring SASL and Schema Registry
For clusters that require SASL authentication or use Schema Registry:
helm install kgazer charts/kgazer \ --set 'kafka.clusters[0].name=production' \ --set 'kafka.clusters[0].bootstrapServers=kafka.example.com:9092' \ --set 'kafka.clusters[0].properties.security\.protocol=SASL_SSL' \ --set 'kafka.clusters[0].properties.sasl\.mechanism=PLAIN' \ --set 'kafka.clusters[0].sasl.username=my-user' \ --set 'kafka.clusters[0].sasl.password=my-password' \ --set 'kafka.clusters[0].schemaRegistry=https://sr.example.com' \ --set 'kafka.clusters[0].schemaRegistryAuth.username=sr-user' \ --set 'kafka.clusters[0].schemaRegistryAuth.password=sr-pass'
SASL credentials and Schema Registry credentials are stored in a Kubernetes Secret rather than the ConfigMap.
Enabling Ingress
helm install kgazer charts/kgazer \ --set ingress.enabled=true \ --set ingress.className=nginx \ --set 'ingress.hosts[0].host=kgazer.example.com' \ --set 'ingress.hosts[0].paths[0].path=/' \ --set 'ingress.hosts[0].paths[0].pathType=Prefix' \ --set 'kafka.clusters[0].name=my-cluster' \ --set 'kafka.clusters[0].bootstrapServers=kafka.example.com:9092'
Helm values reference
| Parameter | Description | Default |
|---|---|---|
replicaCount |
Number of backend replicas | 1 |
image.repository |
Container image | ghcr.io/alfonsojimenez/kgazer |
image.tag |
Image tag (defaults to chart appVersion) |
"" |
kgazer.compactedOnly |
Only consume compacted topics | true |
kafka.clusters |
List of Kafka cluster definitions | [] |
postgresql.enabled |
Deploy bundled PostgreSQL | true |
postgresql.auth.password |
Bundled PostgreSQL password | kgazer |
postgresql.primary.persistence.size |
PostgreSQL storage size | 10Gi |
externalDatabase.enabled |
Use an external PostgreSQL instance | false |
externalDatabase.host |
External database host | "" |
externalDatabase.existingSecret |
Existing Secret for the database password | "" |
ingress.enabled |
Enable Ingress | false |
ingress.className |
Ingress class name | "" |
resources |
CPU/memory requests and limits | {} |
See charts/kgazer/values.yaml for the full list of configurable parameters.
Configuration Reference
| Field | Description | Default |
|---|---|---|
kafka.clusters[].name |
Cluster display name (URL-safe) | required |
kafka.clusters[].bootstrapServers |
Kafka broker addresses | required |
kafka.clusters[].properties |
librdkafka properties (SASL, SSL, etc.) | {} |
kafka.clusters[].schemaRegistry |
Schema Registry URL for Avro deserialisation | — |
kafka.clusters[].schemaRegistryAuth |
Schema Registry credentials (username, password) |
— |
kgazer.server.port |
API server port | 8080 |
kgazer.db.* |
PostgreSQL connection details | — |
kgazer.compactedOnly |
Only consume topics with cleanup.policy=compact |
true |
Environment variable overrides:
| Variable | Overrides |
|---|---|
KGAZER_DB_HOST |
kgazer.db.host |
KGAZER_DB_PASSWORD |
kgazer.db.password |
KGAZER_KAFKA_BOOTSTRAP_SERVERS |
All clusters' bootstrapServers |
How It Works
KGazer has four main components that work together:
┌─────────────────────────────────────────────────────┐
│ Frontend (React) │
│ localhost:4200 │
└──────────────────────┬──────────────────────────────┘
│ REST API
┌──────────────────────▼──────────────────────────────┐
│ Backend (Go) │
│ │
│ ┌──────────┐ ┌──────────┐ ┌───────────────────┐ │
│ │ Syncer │ │ Consumer │ │ API Server │ │
│ │ (per │ │ (per │ │ (chi router) │ │
│ │ cluster) │ │ cluster) │ │ │ │
│ └─────┬─────┘ └─────┬────┘ └────────┬──────────┘ │
│ │ │ │ │
│ ▼ ▼ ▼ │
│ ┌─────────────────────────────────────────────────┐ │
│ │ PostgreSQL │ │
│ │ topics │ messages │ keys │ │
│ └─────────────────────────────────────────────────┘ │
└──────────────────────┬──────────────────────────────┘
│
┌────────▼────────┐
│ Kafka Clusters │
└─────────────────┘
Syncer
One goroutine per cluster, polling Kafka metadata every 60 seconds. On each tick it:
- Calls
GetMetadatato discover all topics in the cluster - Filters to compacted topics only (when
compactedOnlyis enabled) - Calls
DescribeTopicsto get each topic's internal UUID - Upserts topics into PostgreSQL — new topics are inserted, existing ones are updated
- Stale topic cleanup — topics present in the database but missing from metadata are removed (cascade: messages, keys, topic row)
- Recreation detection — if a topic's UUID changed since the last sync, its data is purged and consumption restarts from the beginning
Consumer
One kafka.Consumer per cluster using manual partition assignment (no consumer groups, no rebalancing). It:
- Waits for the initial sync to complete
- Queries the database for all known topics and their stored offsets
- Assigns all partitions across all topics, resuming from
stored_offset + 1(orOffsetBeginningfor new topics) - Polls in a tight loop, batching messages (500 per batch or every 500ms)
- Writes batches to PostgreSQL:
messagestable (withON CONFLICT DO NOTHING),keystable (upsert with latest offset/timestamp), and incremental topic stats - Dynamic topic discovery — every sync interval, compares the database topic list with its in-memory assignment and adds/removes partitions accordingly
- Re-consume support — checks for pending re-consume requests each poll cycle (~100ms) and resets partitions to
OffsetBeginningwhen triggered
Messages are deserialised in this order:
- If Schema Registry is configured and the message has the Confluent wire format (
0x00magic byte) → Avro or Protobuf, depending on the schema's registeredschemaType - If the bytes are valid JSON → JSON
- Otherwise → base64-encoded as
{"_binary": "..."} - Null values (tombstones) → stored as JSON
nullwith formattombstone
API Server
A chi HTTP router exposing a REST API. Key endpoint groups:
| Group | Endpoints | Description |
|---|---|---|
| Topics | GET /api/topics, GET /api/topics/{topic} |
List and detail with consumption progress |
| Keys | GET /api/topics/{topic}/keys |
Paginated key listing with search, partition filter, offset filter and value search |
| Fields | GET /api/topics/{topic}/fields |
Discovered JSON field names for a topic (used by value search autocomplete) |
| History | GET /api/topics/{topic}/keys/{key}/history |
Message versions for a key (most recent first) |
| Timeline | GET /api/topics/{topic}/timeline |
Lightweight version metadata for the change timeline visualisation |
| Consumer Groups | GET /api/topics/{topic}/consumer-groups, GET /{group} |
Live consumer groups with per-partition lag |
| Offset Reset | POST /api/topics/{topic}/consumer-groups/{group}/reset-offsets |
Reset offsets to earliest, latest, or specific values |
| Settings | GET /api/settings/info, GET /api/settings/orphaned-clusters, DELETE /api/settings/clusters/{cluster} |
Runtime stats, sanitised config, orphaned cluster management |
Full API documentation is available in backend/openapi.yaml (OpenAPI 3.0).
Progress Tracking
Consumption progress is tracked in-memory by a dedicated component that syncs Kafka watermarks (low/high offsets) every 30 seconds. Progress percentage is calculated as:
progress = (consumed_offset - low_watermark) / (high_watermark - low_watermark) × 100
A topic is marked as "done" when progress reaches 99.5% or when no new messages have been seen for 30 seconds (indicating the consumer has caught up).
Database Schema
Three tables, defined in backend/migrations/:
topics — one row per cluster+topic combination. Stores partition count, compacted flag, message format, Kafka topic UUID and cached stats (message count, key count, last message timestamp).
messages — every consumed message. Primary key is the composite (topic_id, partition, offset_id) with ON CONFLICT DO NOTHING for idempotent writes. Stores the deserialised body as JSONB.
keys — one row per unique key per topic. Tracks the latest partition, offset, timestamp and a cached copy of the latest message body (JSONB). Used for the key browser with indexes for offset-based sorting, time-based sorting, trigram key search and GIN jsonb_path_ops value search.
Project Structure
kgazer/
├── .github/workflows/ # CI + Release pipelines
├── backend/
│ ├── cmd/server/ # Entry point
│ ├── internal/
│ │ ├── api/ # HTTP handlers (chi router)
│ │ ├── config/ # YAML config loader with validation
│ │ ├── consumer/ # Kafka consumer with dynamic topic assignment
│ │ ├── decoder/ # Avro & Protobuf deserialisation via Schema Registry
│ │ ├── progress/ # Consumption progress tracking
│ │ ├── status/ # Cluster connection status
│ │ ├── store/ # PostgreSQL data access layer
│ │ └── syncer/ # Topic discovery and lifecycle management
│ ├── migrations/ # PostgreSQL schema migrations
│ └── openapi.yaml # API specification (OpenAPI 3.0)
├── frontend/
│ ├── src/
│ │ ├── components/ # Shared UI (JsonViewer, MessageDiff, Sidebar, etc.)
│ │ ├── lib/ # API client and utilities
│ │ ├── pages/ # Route pages (Topics, Keys, History, ConsumerGroups, Settings)
│ │ └── test/ # Test setup
│ └── public/ # Static assets (logo)
├── charts/kgazer/ # Helm chart for Kubernetes deployment
├── docs/ # Documentation assets
├── config.yml.example # Configuration template
├── docker-compose.yml # Development stack
├── VERSION # SemVer version source of truth
└── LICENCE # MIT
URL Structure
KGazer uses path-based cluster routing:
/clusters/:cluster → Topic list
/clusters/:cluster/topics/:topic/keys → Key browser
/clusters/:cluster/topics/:topic/keys/:key → Message history + diffs
/clusters/:cluster/topics/:topic/consumer-groups → Consumer group lag
/settings → Runtime info, config, orphaned clusters
Versioning
KGazer follows Semantic Versioning. The version is defined in the VERSION file at the repository root and injected into the backend binary at build time via -ldflags.
To create a release:
git tag v0.1.0 git push origin v0.1.0
This triggers the GitHub Actions release workflow, which builds and publishes a Docker image to ghcr.io/alfonsojimenez/kgazer:0.1.0.
Development
Running locally (without Docker)
Backend:
cd backend
go run ./cmd/serverRequires Go 1.25+, librdkafka, and a running PostgreSQL instance.
Frontend:
cd frontend
npm install
npm run devThe Vite dev server proxies /api requests to http://localhost:8080. When using Docker Compose, the proxy targets the internal Docker network automatically.
Building for production
cd backend CGO_ENABLED=1 go build -tags musl -o kgazer ./cmd/server cd frontend npm run build # outputs to dist/
CI/CD
Two GitHub Actions workflows in .github/workflows/:
ci.yml— Runs on push/PR tomain. Backend: Go build + vet + tests (with Postgres service). Frontend: TypeScript check + Vitest.release.yml— Runs on tag push (v*). Builds the Docker image with the version baked in and publishes to GitHub Container Registry.
Licence
MIT — see LICENCE.
Trademarks
Apache Kafka and Kafka are either registered trademarks or trademarks of the Apache Software Foundation in the United States and/or other countries. KGazer is not affiliated with, endorsed by, or sponsored by the Apache Software Foundation.
