This project is an Event-Driven Streaming Platform designed to handle high-throughput media events with low latency. It demonstrates how a modern distributed backend system can achieve eventual consistency, scalable search indexing, and real-time analytics using Kafka, Debezium, and Apache Flink.
Built as my Graduation Thesis (Đồ Án Tốt Nghiệp), this repository focuses on solving the challenges of distributed data synchronization and reliable stream processing at scale.
- High Throughput: Capable of ingesting and processing massive volumes of events concurrently.
- Low Latency: Sub-millisecond read access via caching and highly optimized search indexing.
- Fault Tolerance: No single point of failure; robust recovery mechanisms for message processing.
- Scalability: Horizontal scaling capabilities for both stateless microservices and stateful data pipelines.
- Reliability: Zero data loss guarantee via Dead Letter Queues (DLQ) and Change Data Capture (CDC).
- Event-Driven Architecture (EDA): Fully decoupled services reacting to state changes asynchronously.
- Database per Service: Strict data isolation; each microservice manages its own domain database.
- Eventual Consistency: Ensuring data perfectly syncs across services, caches, and search indexes over time.
- CQRS Read Model: Separating heavy write operations (Postgres) from blazing-fast read/search queries (Elasticsearch).
- Loose Coupling: Services communicate via Kafka events rather than direct synchronous HTTP calls where possible.
- Stateless Services: Easy horizontal scaling by keeping session state in Redis.
- API Gateway Pattern: A single unified entry point for all clients, handling routing and auth.
- Change Data Capture (CDC): Capturing row-level database changes directly from the WAL to guarantee event delivery.
code/
├── services/
│ ├── api-gateway # Routing, Load Balancing, JWT Auth verification
│ ├── user-service # Profile & Authentication management
│ ├── content-service # Media metadata and file uploads (MinIO)
│ ├── search-service # Elasticsearch integration for fast querying
│ ├── analytics-service # Viewing and engagement metrics
│ ├── apache-flink-service # Real-time stateful stream processing
│ └── dlq-worker # Dead Letter Queue error recovery
│
├── shared-libraries/
│ └── event-models # Shared Kafka Event Schemas & DTOs
│
├── contracts/
│ └── openapi # OpenAPI/Swagger specifications
│
└── infra/
├── docker-compose # Local development environment (Kafka, Postgres, etc.)
└── kubernetes # Production-ready K8s manifests & Helm charts
The system utilizes an API Gateway pattern to route external traffic to internal, specialized microservices. State changes are captured and streamed via Kafka to keep specialized read models (Elasticsearch) and analytics engines (Flink) updated.
graph TD
Client[Client Applications] -->|REST / HTTP| Gateway[API Gateway]
Gateway --> UserService[User Service]
Gateway --> ContentService[Content Service]
Gateway --> SearchService[Search Service]
UserService -->|Read/Write| UserDB[(PostgreSQL User DB)]
UserService -.->|Cache| RedisUser[(Redis)]
ContentService -->|Read/Write| ContentDB[(PostgreSQL Content DB)]
ContentService -->|Upload/Stream| MinIO[(MinIO Storage)]
ContentDB -->|CDC Stream| Debezium[Debezium Connect]
Debezium -->|Publish Events| Kafka[[Apache Kafka Cluster]]
Kafka -->|Consume Events| SearchService
SearchService -->|Index/Search| Elastic[(Elasticsearch)]
SearchService -.->|Cache| RedisSearch[(Redis)]
Kafka -->|Stream| Flink[Apache Flink]
Flink -->|Aggregated Data| AnalyticsDB[(PostgreSQL Analytics DB)]
Gateway --> AnalyticsService[Analytics Service]
AnalyticsService -->|Read| AnalyticsDB
How a piece of content goes from a user's upload to being searchable in milliseconds without synchronous blocking:
sequenceDiagram
participant User
participant Gateway as API Gateway
participant Content as Content Service
participant DB as PostgreSQL
participant Debezium
participant Kafka
participant Search as Search Service
participant ES as Elasticsearch
User->>Gateway: POST /api/v1/contents
Gateway->>Content: Route Request
Content->>DB: Save Content Metadata (Transaction Commit)
Content-->>User: 201 Created (Fast Response)
DB->>Debezium: Read Write-Ahead Log (WAL)
Debezium->>Kafka: Publish `content.created` Event
Kafka->>Search: Consume Event
Search->>ES: Upsert Document to Index
The backbone of inter-service communication revolves around the following key topics. Partition keys are strictly enforced to guarantee event ordering per entity.
| Topic | Partition Key | Producer | Consumer(s) | Retention | Replication |
|---|---|---|---|---|---|
content.created |
contentId |
Debezium (CDC) | Search Service, Flink | 7 Days | 3 |
content.updated |
contentId |
Debezium (CDC) | Search Service | 7 Days | 3 |
user.registered |
userId |
User Service | Analytics Service | 30 Days | 3 |
content.viewed |
userId |
API Gateway | Flink, Analytics | 3 Days | 3 |
dlq.content |
eventId |
Kafka (Broker) | DLQ Worker | 14 Days | 3 |
Why choose this specific tech stack?
- Why CDC instead of the Outbox Pattern? While the Outbox Pattern solves the dual-write problem by saving events in a database table during the same transaction, it requires application-level logic and constant polling/tailing of that outbox table, increasing database load. Debezium CDC reads directly from the database's Write-Ahead Log (WAL) transparently. This completely decouples event generation from application logic and is significantly more performant.
- Why Kafka instead of RabbitMQ? Kafka provides persistent, append-only logs, high throughput, and consumer replayability. Since we have multiple downstream consumers (Search, Analytics, Flink) needing the same events, Kafka's pub/sub model with consumer groups is far superior to RabbitMQ's transient queues.
- Why CDC (Debezium) instead of Dual Write? Writing to PostgreSQL and Kafka simultaneously in application code (Dual Write) can lead to data inconsistency if the DB commits but the Kafka publish fails (or vice versa). CDC guarantees 100% reliable event emission.
- Why Elasticsearch instead of PostgreSQL Full-Text Search? While PostgreSQL supports text search, Elasticsearch is purpose-built for scalable, distributed full-text indexing, complex aggregations, and fuzzy matching, providing sub-50ms latency over millions of records.
- Why Apache Flink instead of Kafka Streams? Flink offers robust exactly-once state semantics, advanced windowing, and better cluster management for high-volume analytics compared to the library-based approach of Kafka Streams.
Security is integrated at multiple layers of the system to ensure data protection and access control:
- Gateway Authentication: The API Gateway acts as the first line of defense, intercepting all external requests and validating authentication tokens before routing to internal microservices.
- JWT Authentication: Stateless, cryptographically signed JSON Web Tokens (JWT) are used to authenticate user sessions, allowing microservices to verify identity independently without hitting a central database.
- BCrypt Hashing: Passwords are never stored in plain text. They are hashed using the strong, computationally expensive BCrypt algorithm to protect against brute-force attacks.
- Spring Security: Integrated across all Java microservices to secure endpoints and define fine-grained access rules.
- Role-based Authorization (RBAC): Specific endpoints and administrative actions are restricted based on user roles (e.g., User vs. Admin) encoded within the JWT.
- CORS Configuration: Cross-Origin Resource Sharing is strictly configured at the Gateway layer to only allow requests from trusted web domains.
- Payload Validation: Strict input validation using standard Jakarta annotations (e.g.,
@Valid,@NotNull) prevents injection attacks and malformed requests at the controller level.
Building a distributed system introduces complex challenges that were successfully addressed in this project:
- Avoiding Dual Writes: Solved using Change Data Capture (Debezium).
- Data Consistency: Ensured eventual consistency between relational databases and NoSQL search indexes.
- Message Failures & Retries: Implemented Dead Letter Queues (DLQ) and automated retry policies for failed event processing.
- Exactly-Once Processing: Configured Kafka transactions and Flink checkpoints to ensure metrics are not double-counted.
- Handling Backpressure: Kafka acts as a massive buffer to prevent backend services from being overwhelmed during traffic spikes.
The system is designed with the following performance objectives (Load testing verification via JMeter/k6):
- Search Latency:
< 50ms(p95) using Elasticsearch caching. - API Response Time:
< 100msfor core CRUD operations. - Event Processing: High throughput capability supporting
20,000+ events/sec. - End-to-End CDC Latency:
< 500msfrom Postgres transaction commit to Elasticsearch index update.
- Java 17 & Spring Boot 3.2.4
- Spring Cloud Gateway 2023
- Apache Kafka 3.7.0 & Apache Flink 1.19.2
- Debezium 2.7 (CDC)
- PostgreSQL 15 (Primary Relational)
- Elasticsearch 8.13.4 (Search Engine)
- Redis 7 (Distributed Cache)
- MinIO (S3-Compatible Object Storage)
- Stateless Microservices
- JWT Authentication & Bcrypt Hashing
- Kafka Consumer Groups & Transactions
- Flink Stateful Operators & Windowing
- Dockerized Development Environment
- Kubernetes Ready (Helm/Manifests)
- Flyway Database Migrations
api-gateway: Routes traffic, enforces JWT Auth, applies rate limiting.user-service: CRUD Users, Registration, Authentication.content-service: CRUD Content, Uploads media, Saves to DB (Triggers CDC).search-service: Consumes Kafka events, indexes data to ES, serves search queries.analytics-service: Serves aggregated metrics and view counts.apache-flink-service: Real-time stream processing for engagement metrics.dlq-worker: Re-processes failed messages from Kafka Dead Letter Queues.
Here is a glimpse of the primary endpoints exposed via the Gateway:
- Authentication
POST /api/v1/auth/register- Register userPOST /api/v1/auth/login- Authenticate / JWT
- Content
POST /api/v1/contents- Create contentGET /api/v1/contents- List contents
- Search
GET /api/v1/search?q={query}- Full-text search
(Add screenshots of your running system below)
- Swagger API:
<!-- Add Image Link Here --> - Kafka UI:
<!-- Add Image Link Here --> - MinIO Console:
<!-- Add Image Link Here --> - Kibana / Elasticsearch Dashboard:
<!-- Add Image Link Here -->
This graduation thesis is an independent project. I was solely responsible for the end-to-end development:
- Architecture & Design
- Backend Microservices (Spring Boot)
- Data Pipelines & Stream Processing (Kafka, Flink, CDC)
- Database Optimization (PostgreSQL, Elasticsearch)
- DevOps & Infrastructure (Docker, Kubernetes)
- Security (Gateway, JWT)
- Navigate to the local infrastructure directory:
cd code/infra/docker-compose - Start core infrastructure:
docker-compose -f docker-compose.base.yml up -d kafka postgres redis elasticsearch minio debezium kafka-ui
- Start microservices:
docker-compose -f docker-compose.base.yml --profile app up -d
- Access API Gateway at
http://localhost:18081and Kafka UI athttp://localhost:18080.