SyncStream
Distributed data synchronization using PostgreSQL change events, Kafka, Redis, and Elasticsearch.
Tech Stack: Java · PostgreSQL · Debezium · Kafka · Redis · Elasticsearch · Docker
Problem
Synchronizing caches and search indexes through application dual writes creates failure windows and duplicated integration logic.
What I built
Built a CDC pipeline with Debezium, Kafka consumers for Redis and Elasticsearch, shared retries, dead-letter queues, and consumer-management APIs.
Key decisions
CDC through Debezium and Kafka: Change capture separates source writes from downstream projection, at the cost of more infrastructure. Separate cache and search consumers: Consumers can evolve independently while sharing the event stream.
Results
Implemented product-event propagation from PostgreSQL to cache and search projections, with repository verification scripts.
Limitations
A local development platform; analytics consumption is not implemented and role-header authorization needs production hardening.
Evidence
Source, Docker Compose setup, architecture documentation, and workflow verification scripts are available in the repository.
1. Problem Statement
Synchronizing caches and search indexes through application dual writes creates failure windows and duplicated integration logic.
2. Real-World Motivation
Keep PostgreSQL as the source of truth and propagate its changes to independent downstream consumers.
3. System Architecture
Core product-event pipeline at commit 2c5f96e. Both Java consumers use shared bounded retry logic. Supporting services outside this view include the registry and admin UI, SQLite management metadata, Kafka exporter, Prometheus and Grafana. Replay requests are queued as metadata; automatic consumer replay is not implied. Infrastructure is a local development setup.
4. Pipeline Data Flow
Current behavior at commit 2c5f96e: terminal processing errors trigger a dead-letter publish attempt, then offsets are committed. Dead-letter publication errors are logged and swallowed, so durable failure retention and exactly-once processing are not established by this implementation.
5. Failure Modes & Mitigations
| Scenario | Impact | Mitigation Strategy |
|---|---|---|
| Downstream projection fails | Cache or search updates can lag behind PostgreSQL. | Bounded retries are followed by a dead-letter publish attempt. Publish errors are logged; durable failure retention is not guaranteed. |
| Consumer needs recovery | Previously read events may need reprocessing. | Registry APIs persist replay requests for management workflows; consumer replay execution is not wired to those requests. |
6. Design Tradeoffs
| Decision | Alternative | Rationale |
|---|---|---|
| CDC through Debezium and Kafka | Application dual writes or polling | Change capture separates source writes from downstream projection, at the cost of more infrastructure. |
| Separate cache and search consumers | One combined projection service | Consumers can evolve independently while sharing the event stream. |
7. Validation
Deterministic verification scripts cover the monitoring and admin dashboard workflows.
8. Setup & Delivery
Docker Compose starts infrastructure; Maven and Python commands start consumers and the registry platform.
Results & Evaluation
- Product events flow to Redis and Elasticsearch.
- Verification scripts cover monitoring and administration workflows.
Scaling Strategy
Kafka decouples capture from projection. Partitioning and consumer groups provide a path to parallel processing; no production capacity figure is claimed.
Security Model
The local setup uses role headers. Signed identity and production credential management remain deployment work.
Observability
- Prometheus and GrafanaRepository configurations support operational visibility.
- Consumer dashboardRegistration, health, and replay management APIs.
Future Roadmap
Documented next steps include an analytics consumer and production identity controls.