MB_CORE_LOG
Distributed SystemsStatus: Prototype

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

PostgreSQL remains the source of truth; Kafka decouples Redis and Elasticsearch consumers from change capture.

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

The same processing pattern runs independently in the Redis and Elasticsearch consumers.

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

ScenarioImpactMitigation Strategy
Downstream projection failsCache 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 recoveryPreviously 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

DecisionAlternativeRationale
CDC through Debezium and KafkaApplication dual writes or pollingChange capture separates source writes from downstream projection, at the cost of more infrastructure.
Separate cache and search consumersOne combined projection serviceConsumers 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

Evaluation Summary
Functional CDC propagation is implemented; no measured throughput or latency benchmark is published here.
Evaluation Scope
  • 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 Grafana
    Repository configurations support operational visibility.
  • Consumer dashboard
    Registration, health, and replay management APIs.

Future Roadmap

Documented next steps include an analytics consumer and production identity controls.