ShadowStream is a Change Data Capture (CDC) system that streams database changes in real time, buffers them through Redis Streams, archives them to Kafka, and provides controlled replay over gRPC APIs. Built for reliability, scalability, and observability.
Real-time Change Capture — PostgreSQL logical replication via WAL streaming
Dual Storage Pipeline — Redis streams (low latency) + Kafka (durable archival)
Parallel Processing — Kafka consumer groups for high-throughput workloads
Controlled Replay — Speed-controlled, time-travel replay via gRPC
Exactly-Once Semantics — Transactional guarantees across Redis + Kafka
Real-time Monitoring — gRPC progress streaming
Admin Dashboard — Django-based UI for managing replays
Prerequisites
Docker & Docker Compose
Python 3.12+
Clone & Setup git clone https://github.com/mrinalxdev/shadowstream.git cd shadowstream
Start Infrastructure docker-compose up -d postgres redis kafka
Initialize Database docker-compose exec postgres psql -U postgres -d shadowdb
-c "CREATE USER repluser WITH REPLICATION LOGIN PASSWORD 'replpass';"Start Services
# Start ingestor
docker-compose up -d ingestor
# Start control panel (Django)
docker-compose up -d control
# Start replayer
docker-compose up -d replayer- Access Dashboard
Admin Panel: http://localhost:8000/admin
Kafka UI: http://localhost:8081
Environment Variables
PG_HOST=postgres
PG_USER=repluser
PG_PASSWORD=replpass REDIS_URL=redis://redis:6379/0 KAFKA_BROKERS=kafka:9092
shadowstream.archive — archived change events
Consumer groups: analytics-group, backup-group
gRPC Services
serviceIngestor {
rpcPushChange(ChangeRecord) returns (PushResponse);
}
serviceReplayer {
rpcReplayEvent(ReplayRequest) returns (ReplayResponse);
}
serviceProgressTracker {
rpcStreamProgress(ReplayRequest) returns (streamProgressUpdate);
}