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.
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
- Java 21 (Flink 1.20 requires Java 21 or lower)
- Apache Kafka (download from kafka.apache.org)
- Maven 3.8+
# 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.propertiesbin/kafka-topics.sh --create \
--topic log-events \
--partitions 5 \
--replication-factor 1 \
--bootstrap-server localhost:9092cd log-stream-processor/log-producer
mvn spring-boot:run
# Starts on port 8080
# LogSimulator begins generating events every 500ms automaticallycd log-stream-processor/log-processor
mvn exec:java -Dexec.mainClass="com.bishwa.log.processor.FlinkProcessorMain"# 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"}'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
| 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 |
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)