Skip to content

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Latest commit

 

History

159 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 

Repository files navigation

High-Throughput Event-Driven Streaming Platform

Java Spring Boot Apache Kafka Apache Flink PostgreSQL Elasticsearch Docker Kubernetes

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.


Non-functional Goals

  • 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).

Architecture Principles

  • 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.

Project Structure

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

System Architecture

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
Loading

Event Flow: Content Creation to Search Indexing

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
Loading

Kafka Event Topics

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

Design Decisions

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

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.

Challenges Solved

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.

Performance Goals

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: < 100ms for core CRUD operations.
  • Event Processing: High throughput capability supporting 20,000+ events/sec.
  • End-to-End CDC Latency: < 500ms from Postgres transaction commit to Elasticsearch index update.

Technologies & Highlights

Core Stack

  • 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)

Databases & Storage

  • PostgreSQL 15 (Primary Relational)
  • Elasticsearch 8.13.4 (Search Engine)
  • Redis 7 (Distributed Cache)
  • MinIO (S3-Compatible Object Storage)

Technical Highlights

  • Stateless Microservices
  • JWT Authentication & Bcrypt Hashing
  • Kafka Consumer Groups & Transactions
  • Flink Stateful Operators & Windowing
  • Dockerized Development Environment
  • Kubernetes Ready (Helm/Manifests)
  • Flyway Database Migrations

Microservices Overview

  • 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.

Core APIs

Here is a glimpse of the primary endpoints exposed via the Gateway:

  • Authentication
    • POST /api/v1/auth/register - Register user
    • POST /api/v1/auth/login - Authenticate / JWT
  • Content
    • POST /api/v1/contents - Create content
    • GET /api/v1/contents - List contents
  • Search
    • GET /api/v1/search?q={query} - Full-text search

System Screenshots

(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 -->

Project Responsibilities

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)

How to Run Locally

  1. Navigate to the local infrastructure directory:
    cd code/infra/docker-compose
  2. Start core infrastructure:
    docker-compose -f docker-compose.base.yml up -d kafka postgres redis elasticsearch minio debezium kafka-ui
  3. Start microservices:
    docker-compose -f docker-compose.base.yml --profile app up -d
  4. Access API Gateway at http://localhost:18081 and Kafka UI at http://localhost:18080.

About

No description, website, or topics provided.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages