Skip to content

Latest commit

 

History

History
126 lines (102 loc) · 4.36 KB

File metadata and controls

126 lines (102 loc) · 4.36 KB

Log Stream Processor — Kafka + Apache Flink

A real-time log processing pipeline using Spring Boot (producer) and Apache Flink (stream processor).
Built for interview preparation — every class has detailed comments explaining the WHY behind each design decision.

Architecture

LogSimulator (@Scheduled)
      ↓ (every 500ms)
LogEventProducer
      ↓ (KafkaTemplate)
Kafka Topic: log-events (5 partitions)
      ↓ (KafkaSource)
Flink Pipeline:
  filter(ERROR + WARN)
      ↓
  keyBy(service + level)
      ↓
  TumblingEventTimeWindow(60s)
      ↓
  LogLevelCounter (aggregate)  ←— incremental count + sum
      ↓
  ErrorSpikeDetector (process) ←— add window timestamp + spike alert
      ↓
  print() to console           ←— replace with Elasticsearch in prod

Prerequisites

  • Java 21 (Flink 1.20 requires Java 21 or lower)
  • Apache Kafka (download from kafka.apache.org)
  • Maven 3.8+

Quick Start

1. Start Kafka (no Docker needed)

# Terminal 1 — ZooKeeper
cd ~/kafka  # wherever you extracted Kafka
bin/zookeeper-server-start.sh config/zookeeper.properties

# Terminal 2 — Kafka broker
bin/kafka-server-start.sh config/server.properties

2. Create the Kafka topic

bin/kafka-topics.sh --create \
  --topic log-events \
  --partitions 5 \
  --replication-factor 1 \
  --bootstrap-server localhost:9092

3. Start the log producer (Spring Boot)

cd log-stream-processor/log-producer
mvn spring-boot:run
# Starts on port 8080
# LogSimulator begins generating events every 500ms automatically

4. Start the Flink processor

cd log-stream-processor/log-processor
mvn exec:java -Dexec.mainClass="com.bishwa.log.processor.FlinkProcessorMain"

5. (Optional) Inject an error burst via REST

# Trigger 20 ERROR events immediately (makes Flink's spike detector fire)
curl -X POST "http://localhost:8080/api/logs/burst/error?service=payment-service&count=20"

# Send a custom event
curl -X POST http://localhost:8080/api/logs/send \
  -H "Content-Type: application/json" \
  -d '{"level":"ERROR","service":"payment-service","message":"Payment gateway timeout"}'

Expected Output (Flink console after 60s)

LOG STATS> [2025-01-15T10:01:00Z] service=payment-service      level=ERROR count= 12 avgResponseTime=  3241ms
LOG STATS> [2025-01-15T10:01:00Z] service=order-service        level=WARN  count=  8 avgResponseTime=   312ms
LOG STATS> [2025-01-15T10:01:00Z] service=inventory-service    level=ERROR count=  5 avgResponseTime=  2890ms
🚨 ERROR SPIKE DETECTED — service=payment-service errors=12 avgResponseTime=3241ms window=60s

Key Interview Talking Points

Concept Implementation
Kafka partitioning Key = service name → ordering guarantee per service
ACKS="all" Strongest durability — all replicas must acknowledge
Event time windowing Uses LogEvent.timestamp, not Kafka ingestion time
Watermarks 5-second lag tolerance for out-of-order events
AggregateFunction Incremental — doesn't buffer all events in memory
Checkpointing 30s intervals — enables exactly-once on failure/restart
Backpressure Flink slows Kafka reads if downstream is overloaded

Module Structure

log-stream-processor/
├── pom.xml                     # Parent POM (Java 21, 2 modules)
├── log-producer/               # Spring Boot — generates fake logs → Kafka
│   ├── model/LogEvent.java     # Core data model
│   ├── config/KafkaConfig.java # Producer factory, ACKS, serializer
│   ├── service/LogEventProducer.java  # Async Kafka send
│   ├── simulator/LogSimulator.java    # @Scheduled event generator
│   └── controller/LogController.java  # REST API for manual injection
└── log-processor/              # Apache Flink — reads Kafka, aggregates
    ├── FlinkProcessorMain.java        # Entry point, env setup
    ├── job/LogProcessingJob.java      # Full pipeline DAG
    ├── model/LogEvent.java            # Kafka message model (Serializable)
    ├── model/LogStats.java            # Window output model
    ├── deserialization/LogEventDeserializer.java  # bytes → LogEvent
    ├── function/LogLevelCounter.java  # AggregateFunction (count + sum)
    └── function/ErrorSpikeDetector.java  # ProcessWindowFunction (spike alert)