Educational streaming data processing project in Java. The repository combines several lab scenarios with Kafka, Spark Structured Streaming and Apache Flink: event generation, Avro serialization, Kafka ingestion, Parquet output, event-time windows, late/out-of-order events and stateful processing.
| Area | Tools |
|---|---|
| Language | Java 17 |
| Streaming | Apache Kafka, Apache Spark Structured Streaming, Apache Flink |
| Data formats | Avro, Parquet, CSV |
| Build | Maven, Maven Shade Plugin |
| Runtime | Docker Compose, Kafka UI |
- Kafka producer for sending CSV/Avro events to Kafka topics.
- Spark job for reading Avro messages from Kafka and saving results to Parquet.
- Flink job for streaming Avro events to Parquet with savepoint/restart flow.
- Event-time window aggregation in Flink for normal, out-of-order and late events.
- Stateful Flink application with per-user state, TTL and dynamic blocking rules.
src/main/java/
├── lab1/ # Kafka + Spark Structured Streaming + Avro -> Parquet
├── lab2/ # Kafka + Flink + Avro -> Parquet, savepoints
├── lab3/ # Flink event-time windows, late and out-of-order events
└── lab4/ # Flink stateful processing with dynamic blocking rules
Requirements:
- JDK 17
- Maven 3.x
- Docker / Docker Compose
- Apache Spark 3.5.x for lab1
- Apache Flink 2.0.x for lab2-lab4
Build:
mvn -DskipTests clean packageStart Kafka for a lab:
docker compose -f src/main/java/lab3/docker-compose.yml up -dKafka UI is available at http://localhost:8080.
Run a producer for out-of-order events:
mvn -Dexec.mainClass=lab3.producer.EventKafkaProducer \
-Dexec.args="out-of-order" \
exec:javaRun the Flink event-time window job:
$FLINK_HOME/bin/flink run \
-d \
-c lab3.flink.EventTimeWindowCountFlink \
target/sdp-streaming-1.0-SNAPSHOT-all.jar \
user_events_out_of_orderThis project demonstrates local streaming pipeline development, event-time processing, late event handling and stateful stream processing with the Java data engineering stack.