Skip to content

Folders and files

NameName
Last commit message
Last commit date

Latest commit

 

History

7 Commits
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

Streaming Data Processing

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.

Tech Stack

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

Implemented Scenarios

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

Repository Structure

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

Quick Start

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 package

Start Kafka for a lab:

docker compose -f src/main/java/lab3/docker-compose.yml up -d

Kafka UI is available at http://localhost:8080.

Example Commands

Run a producer for out-of-order events:

mvn -Dexec.mainClass=lab3.producer.EventKafkaProducer \
  -Dexec.args="out-of-order" \
  exec:java

Run 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_order

Portfolio Notes

This project demonstrates local streaming pipeline development, event-time processing, late event handling and stateful stream processing with the Java data engineering stack.

About

Streaming data processing labs with Kafka, Spark Structured Streaming and Apache Flink in Java.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages