Skip to content

Repository files navigation

Kafka-Flink-Iceberg

A multi-module Java project that streams financial trading orders (FIX NewOrderSingle) through Kafka, processes them with Apache Flink, and persists to Apache Iceberg tables in S3-compatible storage (MinIO).

Architecture

                                    Podman (kafka-flink-net)
┌──────────────┐    ┌──────────────┐    ┌──────────────┐    ┌──────────────┐
│   Producer   │───>│    Kafka     │───>│    Flink     │───>│    MinIO     │
│  (protobuf)  │    │   (KRaft)   │    │  (streaming) │    │  (Iceberg)   │
└──────────────┘    └──────────────┘    └──────────────┘    └──────────────┘
                          │                                       │
                     ┌────┴─────┐                           ┌─────┴──────┐
                     │ Kafka UI │                           │   DuckDB   │
                     │ :8080    │                           │  (queries) │
                     └──────────┘                           └────────────┘

Modules

Module Description
protobuf NewOrderSingle protobuf definition and generated Java classes
kafka Kafka producer/consumer and Podman scripts for Kafka + Kafka UI
flink Flink streaming job, Iceberg sink, and Podman scripts for Flink + MinIO
duckdb Scripts to query Iceberg tables using DuckDB

Prerequisites

  • Java 17
  • Podman with a running machine (podman machine start)
  • DuckDB CLI (for querying, optional)

Quick Start

1. Build all modules

./gradlew build

2. Start Kafka

# First time only — pull images, create network, generate cluster ID
./kafka/podman/setup-kafka.sh

# Start Kafka broker + Kafka UI
./kafka/podman/start-kafka.sh

3. Start Flink & MinIO

# First time only — pull images, create volumes
./flink/podman/setup-flink.sh

# Start MinIO, Flink JobManager, and TaskManager
./flink/podman/start-flink.sh

4. Submit the Flink job

./flink/podman/submit-job.sh

This builds the shadow JAR, copies it to the Flink JobManager, and submits the NewOrderSingle streaming job.

5. Send orders to Kafka

./gradlew :kafka:runProducer

Sends 100 random NewOrderSingle orders to the new-order-single topic.

6. Query data with DuckDB

# First time only
./duckdb/scripts/setup-duckdb.sh

# Run queries against Iceberg table in MinIO
./duckdb/scripts/query-orders.sh

Podman Containers

Container Image Ports Purpose
kafka-broker apache/kafka:3.9.0 9092 Kafka broker (KRaft mode)
kafka-ui provectuslabs/kafka-ui 8080 Kafka management UI
minio minio/minio 9000, 9001 S3-compatible storage for Iceberg
flink-jobmanager flink:1.20.1-java17 8081 Flink job orchestration
flink-taskmanager flink:1.20.1-java17 Flink task execution

All containers run on the shared kafka-flink-net Podman network.

Managing containers

# Check status
podman ps

# Stop all
./flink/podman/stop-flink.sh
podman stop kafka-broker kafka-ui

# Restart all
podman start kafka-broker kafka-ui minio flink-jobmanager flink-taskmanager

Flink Job Management

# List running jobs
podman exec flink-jobmanager flink list

# Cancel a job
podman exec flink-jobmanager flink cancel <job-id>

# View logs
podman logs flink-jobmanager
podman logs flink-taskmanager

Kafka Management

# List topics
podman exec kafka-broker /opt/kafka/bin/kafka-topics.sh --list --bootstrap-server kafka-broker:19092

# Consume messages (CLI)
./gradlew :kafka:runConsumer

Or use the Kafka UI at http://localhost:8080.

Project Structure

├── build.gradle              # Common Java 17 config
├── settings.gradle           # Module includes + plugin management
├── gradle.properties         # Centralized dependency versions
├── protobuf/
│   ├── build.gradle
│   └── src/main/proto/NewOrderSingle.proto
├── kafka/
│   ├── build.gradle
│   ├── podman/
│   │   ├── setup-kafka.sh    # Pull images, create network
│   │   └── start-kafka.sh    # Start Kafka + Kafka UI
│   └── src/main/java/com/example/kafka/
│       ├── NewOrderSingleSerializer.java
│       ├── NewOrderSingleProducer.java
│       └── NewOrderSingleConsumer.java
├── flink/
│   ├── build.gradle
│   ├── podman/
│   │   ├── setup-flink.sh    # Pull images, create volumes
│   │   ├── start-flink.sh    # Start MinIO + Flink cluster
│   │   ├── stop-flink.sh     # Stop all Flink services
│   │   └── submit-job.sh     # Build + deploy + submit job
│   └── src/main/java/com/example/flink/
│       ├── NewOrderSingleFlinkJob.java
│       ├── NewOrderSingleDeserializationSchema.java
│       └── IcebergCatalogConfig.java
└── duckdb/
    ├── README.md
    └── scripts/
        ├── setup-duckdb.sh   # Install DuckDB + extensions
        ├── query-orders.sh   # Run queries
        └── query-orders.sql  # Default SQL queries

Key Versions

Component Version
Java 17
Protobuf 4.29.3
Kafka 3.9.0 (KRaft)
Flink 1.20.1
Iceberg 1.7.1
Hadoop 3.4.1

About

kafka to apache flink iceberg

Resources

Stars

Watchers

Forks

Releases

Packages

Contributors

Languages