MyCodeSchool
Apache Kafka cu Spring Boot
mycodeschool.ro
+

Apache Kafka + Spring Boot

Tutorial complet de inițiere în Apache Kafka și integrarea cu Spring Boot pentru aplicații distribuite și procesare de evenimente în timp real

Introducere în Apache Kafka

Apache Kafka este o platformă distribuită de streaming de evenimente, concepută pentru a gestiona fluxuri mari de date în timp real. Dezvoltat inițial de LinkedIn și acum parte a ecosistemului Apache, Kafka a devenit standardul industrial pentru construirea de sisteme de mesagerie scalabile și tolerante la erori.

📚 Ce este Kafka?

Kafka este un sistem distribuit de mesagerie de tip publish-subscribe, optimizat pentru throughput ridicat, latență scăzută și durabilitate. Este folosit de companii precum Netflix, Uber, LinkedIn și Airbnb pentru procesarea a miliarde de evenimente zilnic.

De ce să folosești Kafka?

Performanță Ridicată
Poate procesa milioane de mesaje pe secundă cu latență sub-milisecundă
📈
Scalabilitate Orizontală
Adaugă mai multe noduri pentru a crește capacitatea fără downtime
💾
Persistență
Mesajele sunt stocate pe disk și replicate pentru durabilitate
🔄
Toleranță la Erori
Replicare automată și failover pentru disponibilitate ridicată

Arhitectura de Bază

Producer
Trimite mesaje
Kafka Broker
Topic → Partitions
Consumer
Citește mesaje

Concepte Fundamentale

Înainte de a începe implementarea, este esențial să înțelegi terminologia și conceptele de bază ale Kafka.

Concept Descriere
Broker Un server Kafka care stochează date și servește clienți. Un cluster Kafka conține mai mulți brokeri.
Topic O categorie sau un feed la care sunt publicate mesajele. Similar cu un tabel într-o bază de date.
Partition O subdiviziune a unui topic. Permite paralelizarea și scalabilitatea.
Producer Aplicația care publică (scrie) mesaje către un topic Kafka.
Consumer Aplicația care se abonează și citește mesaje dintr-un topic.
Consumer Group Un grup de consumeri care împart procesarea partițiilor unui topic.
Offset Un identificator unic pentru fiecare mesaj într-o partiție, folosit pentru tracking.
Zookeeper/KRaft Serviciu pentru coordonarea cluster-ului. KRaft este noua alternativă fără Zookeeper.
💡 Sfat Important

Numărul de partiții determină gradul maxim de paralelism pentru consumeri. Dacă ai 4 partiții, poți avea maximum 4 consumeri activi în același consumer group.

Fluxul Mesajelor

  1. Producer trimite un mesaj către un Topic
  2. Kafka determină Partition-ul (bazat pe key sau round-robin)
  3. Mesajul primește un Offset unic și este stocat
  4. Consumer citește mesajul și actualizează offset-ul procesat

Instalare și Setup

Opțiunea 1: Docker (Recomandat)

Cea mai rapidă metodă pentru dezvoltare locală este folosind Docker Compose:

docker-compose.yml
version: '3.8'

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    container_name: zookeeper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    ports:
      - "2181:2181"

  kafka:
    image: confluentinc/cp-kafka:7.5.0
    container_name: kafka
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
Terminal
# Pornește Kafka și Zookeeper
docker-compose up -d

# Verifică că rulează
docker-compose ps

# Vezi log-urile
docker-compose logs -f kafka

Opțiunea 2: KRaft Mode (Fără Zookeeper)

Kafka 3.3+ suportă modul KRaft, care elimină dependența de Zookeeper:

docker-compose-kraft.yml
version: '3.8'

services:
  kafka:
    image: confluentinc/cp-kafka:7.5.0
    container_name: kafka-kraft
    ports:
      - "9092:9092"
    environment:
      KAFKA_NODE_ID: 1
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: 'CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT'
      KAFKA_ADVERTISED_LISTENERS: 'PLAINTEXT://localhost:9092'
      KAFKA_PROCESS_ROLES: 'broker,controller'
      KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka:29093'
      KAFKA_LISTENERS: 'PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:29093'
      KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER'
      CLUSTER_ID: 'MkU3OEVBNTcwNTJENDM2Qk'

Crearea Proiectului Spring Boot

1
Generează proiectul

Accesează start.spring.io sau folosește Spring Initializr din IDE-ul tău preferat.

2
Selectează dependențele

Adaugă: Spring for Apache Kafka, Spring Web, Lombok (opțional)

pom.xml
<dependencies>
    <!-- Spring Boot Starter -->
    <dependency>
        <groupId>org.springframework.boot</groupId>
        <artifactId>spring-boot-starter-web</artifactId>
    </dependency>

    <!-- Spring Kafka -->
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka</artifactId>
    </dependency>

    <!-- Lombok (opțional) -->
    <dependency>
        <groupId>org.projectlombok</groupId>
        <artifactId>lombok</artifactId>
        <optional>true</optional>
    </dependency>

    <!-- Jackson pentru JSON -->
    <dependency>
        <groupId>com.fasterxml.jackson.core</groupId>
        <artifactId>jackson-databind</artifactId>
    </dependency>

    <!-- Test -->
    <dependency>
        <groupId>org.springframework.kafka</groupId>
        <artifactId>spring-kafka-test</artifactId>
        <scope>test</scope>
    </dependency>
</dependencies>

Configurare Spring Boot

Configurare de Bază

application.yml
spring:
  kafka:
    # Adresa broker-ului Kafka
    bootstrap-servers: localhost:9092
    
    # Configurare Producer
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      # Garantează că mesajele sunt primite de toți replica
      acks: all
      # Retry în caz de eroare
      retries: 3
    
    # Configurare Consumer
    consumer:
      group-id: my-app-group
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      # Citește de la început dacă nu există offset salvat
      auto-offset-reset: earliest
      properties:
        spring.json.trusted.packages: "*"

# Topic-uri personalizate
app:
  kafka:
    topics:
      orders: orders-topic
      notifications: notifications-topic

Configurare Java (Alternativă)

KafkaConfig.java
package com.example.config;

import org.apache.kafka.clients.admin.NewTopic;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.TopicBuilder;

@Configuration
public class KafkaConfig {

    // Creează topic-ul automat la pornirea aplicației
    @Bean
    public NewTopic ordersTopic() {
        return TopicBuilder.name("orders-topic")
                .partitions(3)           // 3 partiții pentru paralelism
                .replicas(1)              // 1 replică (pentru dev)
                .compact()                 // Compaction policy
                .build();
    }

    @Bean
    public NewTopic notificationsTopic() {
        return TopicBuilder.name("notifications-topic")
                .partitions(2)
                .replicas(1)
                .build();
    }
}
⚠️ Atenție la Producție

În producție, setează replicas la minimum 3 pentru toleranță la erori și acks=all pentru durabilitate maximă.

Implementarea Producer-ului

Model de Date

Order.java
package com.example.model;

import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;

import java.math.BigDecimal;
import java.time.LocalDateTime;

@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class Order {
    
    private String orderId;
    private String customerId;
    private String productName;
    private Integer quantity;
    private BigDecimal totalPrice;
    private OrderStatus status;
    private LocalDateTime createdAt;

    public enum OrderStatus {
        PENDING, CONFIRMED, SHIPPED, DELIVERED, CANCELLED
    }
}

Serviciul Producer

OrderProducerService.java
package com.example.service;

import com.example.model.Order;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.support.SendResult;
import org.springframework.stereotype.Service;

import java.util.concurrent.CompletableFuture;

@Service
@Slf4j
@RequiredArgsConstructor
public class OrderProducerService {

    private final KafkaTemplate<String, Order> kafkaTemplate;

    @Value("${app.kafka.topics.orders}")
    private String ordersTopic;

    /**
     * Trimite o comandă către Kafka (async)
     */
    public CompletableFuture<SendResult<String, Order>> sendOrder(Order order) {
        log.info("📤 Trimit comanda: {}", order.getOrderId());
        
        // Folosim orderId ca key pentru partitioning consistent
        return kafkaTemplate.send(ordersTopic, order.getOrderId(), order)
                .whenComplete((result, ex) -> {
                    if (ex == null) {
                        log.info("✅ Comandă trimisă cu succes: {} -> partition {}, offset {}",
                                order.getOrderId(),
                                result.getRecordMetadata().partition(),
                                result.getRecordMetadata().offset());
                    } else {
                        log.error("❌ Eroare la trimiterea comenzii: {}", 
                                order.getOrderId(), ex);
                    }
                });
    }

    /**
     * Trimite o comandă sincron (blochează până primește confirmare)
     */
    public SendResult<String, Order> sendOrderSync(Order order) {
        try {
            return kafkaTemplate.send(ordersTopic, order.getOrderId(), order)
                    .get(); // Blochează până primește rezultatul
        } catch (Exception e) {
            log.error("❌ Eroare la trimitere sincronă: {}", e.getMessage());
            throw new RuntimeException("Failed to send order", e);
        }
    }
}

Controller REST

OrderController.java
package com.example.controller;

import com.example.model.Order;
import com.example.service.OrderProducerService;
import lombok.RequiredArgsConstructor;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.*;

import java.math.BigDecimal;
import java.time.LocalDateTime;
import java.util.UUID;

@RestController
@RequestMapping("/api/orders")
@RequiredArgsConstructor
public class OrderController {

    private final OrderProducerService producerService;

    @PostMapping
    public ResponseEntity<Order> createOrder(@RequestBody OrderRequest request) {
        Order order = Order.builder()
                .orderId(UUID.randomUUID().toString())
                .customerId(request.customerId())
                .productName(request.productName())
                .quantity(request.quantity())
                .totalPrice(request.price().multiply(BigDecimal.valueOf(request.quantity())))
                .status(Order.OrderStatus.PENDING)
                .createdAt(LocalDateTime.now())
                .build();

        producerService.sendOrder(order);
        
        return ResponseEntity.ok(order);
    }

    // Record pentru request (Java 17+)
    public record OrderRequest(
            String customerId,
            String productName,
            Integer quantity,
            BigDecimal price
    ) {}
}
🍃 Spring Kafka Template

KafkaTemplate este componenta principală pentru trimiterea mesajelor. Spring Boot auto-configurează acest bean bazat pe proprietățile din application.yml.

Implementarea Consumer-ului

Consumer Simplu

OrderConsumerService.java
package com.example.service;

import com.example.model.Order;
import lombok.extern.slf4j.Slf4j;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Service;

@Service
@Slf4j
public class OrderConsumerService {

    /**
     * Consumer simplu - procesează fiecare mesaj
     */
    @KafkaListener(
        topics = "${app.kafka.topics.orders}",
        groupId = "order-processing-group"
    )
    public void consumeOrder(Order order) {
        log.info("📥 Comandă primită: {} - Status: {}", 
                order.getOrderId(), order.getStatus());
        
        // Procesează comanda
        processOrder(order);
    }

    /**
     * Consumer cu acces la metadata (partition, offset, etc.)
     */
    @KafkaListener(
        topics = "${app.kafka.topics.orders}",
        groupId = "order-analytics-group"
    )
    public void consumeWithMetadata(ConsumerRecord<String, Order> record) {
        log.info("📊 Metadata - Topic: {}, Partition: {}, Offset: {}, Key: {}",
                record.topic(),
                record.partition(),
                record.offset(),
                record.key());
        
        Order order = record.value();
        // Procesare pentru analytics...
    }

    private void processOrder(Order order) {
        // Logică de business
        log.info("🔄 Procesez comanda {} pentru clientul {}", 
                order.getOrderId(), order.getCustomerId());
        
        // Simulare procesare
        try {
            Thread.sleep(100); // Simulare delay
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
        
        log.info("✅ Comandă procesată: {}", order.getOrderId());
    }
}

Consumer cu Manual Acknowledgment

ManualAckConsumer.java
package com.example.service;

import com.example.model.Order;
import lombok.extern.slf4j.Slf4j;
import org.springframework.kafka.annotation.KafkaListener;
import org.springframework.kafka.support.Acknowledgment;
import org.springframework.stereotype.Service;

@Service
@Slf4j
public class ManualAckConsumer {

    /**
     * Consumer cu acknowledge manual
     * Necesită: spring.kafka.listener.ack-mode=MANUAL
     */
    @KafkaListener(
        topics = "orders-topic",
        groupId = "manual-ack-group",
        containerFactory = "kafkaListenerContainerFactory"
    )
    public void consumeWithManualAck(Order order, Acknowledgment ack) {
        try {
            log.info("📥 Procesez: {}", order.getOrderId());
            
            // Procesare...
            processOrderSafely(order);
            
            // Confirmă doar după procesare reușită
            ack.acknowledge();
            log.info("✅ Acknowledged: {}", order.getOrderId());
            
        } catch (Exception e) {
            log.error("❌ Eroare procesare: {} - {}", 
                    order.getOrderId(), e.getMessage());
            // NU facem acknowledge - mesajul va fi reprocesat
            // Sau: ack.nack(Duration.ofSeconds(10)); pentru retry cu delay
        }
    }

    private void processOrderSafely(Order order) {
        // Logică de business care poate eșua
        if (order.getTotalPrice().doubleValue() < 0) {
            throw new IllegalArgumentException("Invalid price");
        }
    }
}

Configurare pentru Manual Acknowledgment

KafkaConsumerConfig.java
package com.example.config;

import com.example.model.Order;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.config.ConcurrentKafkaListenerContainerFactory;
import org.springframework.kafka.core.ConsumerFactory;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
import org.springframework.kafka.listener.ContainerProperties;
import org.springframework.kafka.support.serializer.JsonDeserializer;

import java.util.HashMap;
import java.util.Map;

@Configuration
public class KafkaConsumerConfig {

    @Value("${spring.kafka.bootstrap-servers}")
    private String bootstrapServers;

    @Bean
    public ConsumerFactory<String, Order> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class);
        props.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean
    public ConcurrentKafkaListenerContainerFactory<String, Order> 
            kafkaListenerContainerFactory() {
        
        ConcurrentKafkaListenerContainerFactory<String, Order> factory = 
                new ConcurrentKafkaListenerContainerFactory<>();
        
        factory.setConsumerFactory(consumerFactory());
        
        // Activează manual acknowledgment
        factory.getContainerProperties()
                .setAckMode(ContainerProperties.AckMode.MANUAL);
        
        // Număr de thread-uri concurente
        factory.setConcurrency(3);
        
        return factory;
    }
}

Funcționalități Avansate

Batch Consumer

Procesează mai multe mesaje simultan pentru performanță îmbunătățită:

BatchConsumer.java
@KafkaListener(
    topics = "orders-topic",
    groupId = "batch-group",
    containerFactory = "batchFactory"
)
public void consumeBatch(List<Order> orders) {
    log.info("📦 Procesez batch de {} comenzi", orders.size());
    
    orders.forEach(this::processOrder);
    
    log.info("✅ Batch procesat");
}

Error Handling cu Dead Letter Topic

ErrorHandlingConfig.java
package com.example.config;

import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.listener.DeadLetterPublishingRecoverer;
import org.springframework.kafka.listener.DefaultErrorHandler;
import org.springframework.util.backoff.FixedBackOff;

@Configuration
public class ErrorHandlingConfig {

    @Bean
    public DefaultErrorHandler errorHandler(
            KafkaTemplate<String, Object> kafkaTemplate) {
        
        // Publică mesajele eșuate într-un Dead Letter Topic
        DeadLetterPublishingRecoverer recoverer = 
            new DeadLetterPublishingRecoverer(kafkaTemplate,
                (record, ex) -> {
                    // Trimite la topic-ul original + ".DLT"
                    return new TopicPartition(
                        record.topic() + ".DLT", 
                        record.partition()
                    );
                });

        // Retry de 3 ori cu interval de 1 secundă
        return new DefaultErrorHandler(
            recoverer, 
            new FixedBackOff(1000L, 3)
        );
    }
}

Tranzacții Kafka

TransactionalProducer.java
@Service
@RequiredArgsConstructor
public class TransactionalOrderService {

    private final KafkaTemplate<String, Order> kafkaTemplate;

    /**
     * Trimite multiple mesaje într-o singură tranzacție
     * Fie toate reușesc, fie niciuna
     */
    @Transactional
    public void processOrderWithNotification(Order order) {
        // Prima operație
        kafkaTemplate.send("orders-topic", order.getOrderId(), order);
        
        // A doua operație
        Notification notification = createNotification(order);
        kafkaTemplate.send("notifications-topic", 
                order.getCustomerId(), notification);
        
        // Dacă oricare eșuează, ambele sunt rollback
    }
}
💡 Activare Tranzacții

Adaugă în application.yml:

spring.kafka.producer.transaction-id-prefix: tx-

Monitoring cu Actuator

application.yml
management:
  endpoints:
    web:
      exposure:
        include: health,metrics,kafka
  health:
    kafka:
      enabled: true

Best Practices

🔑
Folosește Message Keys
Mesajele cu aceeași key merg întotdeauna în aceeași partiție, garantând ordinea.
📊
Dimensionează Partițiile
Numărul de partiții = numărul maxim de consumeri concurenți. Planifică pentru creștere.
🔄
Idempotență
Proiectează consumerii să fie idempotenți - procesarea repetată a aceluiași mesaj să nu cauzeze probleme.
📝
Schema Registry
Folosește Confluent Schema Registry cu Avro/Protobuf pentru evoluția schemelor.
🛡️
Dead Letter Topics
Trimite mesajele eșuate într-un DLT pentru analiză și reprocessare ulterioară.
📈
Monitorizare
Monitorizează consumer lag, throughput și latență cu Prometheus/Grafana.

Checklist pentru Producție

Kafka UI - Instrumente de Administrare

Pentru a administra și monitoriza clusterele Kafka mai ușor, există mai multe instrumente cu interfață grafică. Acestea permit vizualizarea topic-urilor, mesajelor, consumer groups și multe altele fără a folosi linia de comandă.

Opțiuni Populare pentru Kafka UI

🎯
Kafka UI (Provectus)
Open-source, ușor de configurat, suport pentru multiple clustere. Recomandat pentru dezvoltare și producție.
🔷
Confluent Control Center
Soluție enterprise de la Confluent cu monitoring avansat, alerting și management complet.
📊
AKHQ (ex-KafkaHQ)
Interfață modernă cu suport pentru Kafka Connect, Schema Registry și ACLs.
🦅
Kafdrop
Lightweight, simplu de folosit, ideal pentru development și debugging rapid.

1. Kafka UI (Provectus) - Recomandat

Kafka UI de la Provectus este cel mai popular tool open-source pentru administrarea Kafka. Oferă o interfață modernă și intuitivă.

docker-compose-with-ui.yml
version: '3.8'

services:
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    container_name: zookeeper
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
      ZOOKEEPER_TICK_TIME: 2000
    ports:
      - "2181:2181"
    networks:
      - kafka-network

  kafka:
    image: confluentinc/cp-kafka:7.5.0
    container_name: kafka
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
      - "29092:29092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
    networks:
      - kafka-network

  # ========================================
  # KAFKA UI - Interfață Web pentru Kafka
  # ========================================
  kafka-ui:
    image: provectuslabs/kafka-ui:latest
    container_name: kafka-ui
    depends_on:
      - kafka
    ports:
      - "8080:8080"
    environment:
      # Configurare cluster
      KAFKA_CLUSTERS_0_NAME: local-cluster
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092
      KAFKA_CLUSTERS_0_ZOOKEEPER: zookeeper:2181
      
      # Opțional: Schema Registry (dacă folosești)
      # KAFKA_CLUSTERS_0_SCHEMAREGISTRY: http://schema-registry:8081
      
      # Opțional: Kafka Connect (dacă folosești)
      # KAFKA_CLUSTERS_0_KAFKACONNECT_0_NAME: connect
      # KAFKA_CLUSTERS_0_KAFKACONNECT_0_ADDRESS: http://kafka-connect:8083
      
      # Setări UI
      DYNAMIC_CONFIG_ENABLED: "true"
    networks:
      - kafka-network

networks:
  kafka-network:
    driver: bridge
Terminal
# Pornește toate serviciile
docker-compose -f docker-compose-with-ui.yml up -d

# Verifică că toate containerele rulează
docker-compose -f docker-compose-with-ui.yml ps

# Accesează Kafka UI în browser
# http://localhost:8080
💡 Accesare Kafka UI

După pornire, accesează http://localhost:8080 în browser pentru a vedea interfața Kafka UI.

Funcționalități Kafka UI

Funcționalitate Descriere
📋 Topics Management Crează, șterge, editează topic-uri. Vizualizează partiții, replication factor și configurări.
📨 Message Browser Vizualizează și caută mesaje în topic-uri. Suport pentru JSON, Avro, Protobuf.
👥 Consumer Groups Monitorizează consumer groups, offset-uri și lag. Reset offset-uri manual.
🖥️ Brokers Vizualizează starea brokerilor, configurații și metrici.
📤 Produce Messages Trimite mesaje direct din UI pentru testare rapidă.
📊 Schema Registry Gestionează scheme Avro/Protobuf (dacă e configurat).
🔗 Kafka Connect Administrează conectori pentru integrări externe.
🔍 KSQL Rulează query-uri KSQL pentru stream processing.

2. Configurare pentru Multiple Clustere

Kafka UI suportă monitorizarea mai multor clustere simultan:

Multi-Cluster Config
kafka-ui:
    image: provectuslabs/kafka-ui:latest
    environment:
      # Cluster 1 - Development
      KAFKA_CLUSTERS_0_NAME: development
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka-dev:9092
      KAFKA_CLUSTERS_0_READONLY: "false"
      
      # Cluster 2 - Staging
      KAFKA_CLUSTERS_1_NAME: staging
      KAFKA_CLUSTERS_1_BOOTSTRAPSERVERS: kafka-staging:9092
      KAFKA_CLUSTERS_1_READONLY: "false"
      
      # Cluster 3 - Production (read-only pentru siguranță)
      KAFKA_CLUSTERS_2_NAME: production
      KAFKA_CLUSTERS_2_BOOTSTRAPSERVERS: kafka-prod:9092
      KAFKA_CLUSTERS_2_READONLY: "true"

3. Configurare cu Autentificare

Kafka UI cu Autentificare
kafka-ui:
    image: provectuslabs/kafka-ui:latest
    environment:
      # Activează autentificarea în UI
      AUTH_TYPE: "LOGIN_FORM"
      SPRING_SECURITY_USER_NAME: admin
      SPRING_SECURITY_USER_PASSWORD: admin-secret-password
      
      # Cluster cu SASL autentificare
      KAFKA_CLUSTERS_0_NAME: secured-cluster
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:9092
      KAFKA_CLUSTERS_0_PROPERTIES_SECURITY_PROTOCOL: SASL_PLAINTEXT
      KAFKA_CLUSTERS_0_PROPERTIES_SASL_MECHANISM: PLAIN
      KAFKA_CLUSTERS_0_PROPERTIES_SASL_JAAS_CONFIG: 'org.apache.kafka.common.security.plain.PlainLoginModule required username="kafka" password="kafka-secret";'

4. AKHQ - Alternativă Populară

docker-compose-akhq.yml
version: '3.8'

services:
  akhq:
    image: tchiotludo/akhq:latest
    container_name: akhq
    ports:
      - "8080:8080"
    environment:
      AKHQ_CONFIGURATION: |
        akhq:
          connections:
            local-cluster:
              properties:
                bootstrap.servers: "kafka:29092"
              # Schema Registry (opțional)
              schema-registry:
                url: "http://schema-registry:8081"
              # Kafka Connect (opțional)
              connect:
                - name: "connect"
                  url: "http://kafka-connect:8083"
    networks:
      - kafka-network

5. Kafdrop - Lightweight Option

docker-compose-kafdrop.yml
version: '3.8'

services:
  kafdrop:
    image: obsidiandynamics/kafdrop:latest
    container_name: kafdrop
    ports:
      - "9000:9000"
    environment:
      KAFKA_BROKERCONNECT: kafka:29092
      JVM_OPTS: "-Xms32M -Xmx64M"
      # Pentru Schema Registry
      # SCHEMAREGISTRY_CONNECT: "http://schema-registry:8081"
    networks:
      - kafka-network

6. Integrare cu Spring Boot pentru Metrici

Pentru monitorizare avansată, poți expune metrici din aplicația Spring Boot:

pom.xml - Dependențe
<!-- Actuator pentru metrici -->
<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-actuator</artifactId>
</dependency>

<!-- Micrometer pentru Prometheus -->
<dependency>
    <groupId>io.micrometer</groupId>
    <artifactId>micrometer-registry-prometheus</artifactId>
</dependency>
application.yml - Metrici Kafka
management:
  endpoints:
    web:
      exposure:
        include: health,info,metrics,prometheus
  metrics:
    tags:
      application: ${spring.application.name}
  health:
    kafka:
      enabled: true

# Activează metrici detaliate pentru Kafka
spring:
  kafka:
    producer:
      properties:
        metrics.recording.level: DEBUG
    consumer:
      properties:
        metrics.recording.level: DEBUG

7. Stack Complet cu Monitoring

Un setup complet pentru development cu Kafka, UI, și monitoring:

docker-compose-full-stack.yml
version: '3.8'

services:
  # Zookeeper
  zookeeper:
    image: confluentinc/cp-zookeeper:7.5.0
    environment:
      ZOOKEEPER_CLIENT_PORT: 2181
    ports:
      - "2181:2181"
    networks:
      - kafka-network

  # Kafka Broker
  kafka:
    image: confluentinc/cp-kafka:7.5.0
    depends_on:
      - zookeeper
    ports:
      - "9092:9092"
    environment:
      KAFKA_BROKER_ID: 1
      KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:29092,PLAINTEXT_HOST://localhost:9092
      KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
      KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT
      KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
      KAFKA_JMX_PORT: 9991
    networks:
      - kafka-network

  # Schema Registry
  schema-registry:
    image: confluentinc/cp-schema-registry:7.5.0
    depends_on:
      - kafka
    ports:
      - "8081:8081"
    environment:
      SCHEMA_REGISTRY_HOST_NAME: schema-registry
      SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: kafka:29092
    networks:
      - kafka-network

  # Kafka UI
  kafka-ui:
    image: provectuslabs/kafka-ui:latest
    depends_on:
      - kafka
      - schema-registry
    ports:
      - "8080:8080"
    environment:
      KAFKA_CLUSTERS_0_NAME: local
      KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka:29092
      KAFKA_CLUSTERS_0_ZOOKEEPER: zookeeper:2181
      KAFKA_CLUSTERS_0_SCHEMAREGISTRY: http://schema-registry:8081
      KAFKA_CLUSTERS_0_METRICS_PORT: 9991
      DYNAMIC_CONFIG_ENABLED: "true"
    networks:
      - kafka-network

networks:
  kafka-network:
    driver: bridge
⚠️ Porturi Folosite

2181 - Zookeeper
9092 - Kafka Broker
8081 - Schema Registry
8080 - Kafka UI

Comparație Instrumente UI

Caracteristică Kafka UI AKHQ Kafdrop Confluent CC
Licență Open Source Open Source Open Source Commercial
Browse Messages ✅ Da ✅ Da ✅ Da ✅ Da
Produce Messages ✅ Da ✅ Da ❌ Nu ✅ Da
Schema Registry ✅ Da ✅ Da ✅ Da ✅ Da
Kafka Connect ✅ Da ✅ Da ❌ Nu ✅ Da
ACL Management ✅ Da ✅ Da ❌ Nu ✅ Da
Multi-Cluster ✅ Da ✅ Da ❌ Nu ✅ Da
Memorie RAM ~200MB ~300MB ~100MB >1GB
Recomandare Dev & Prod Dev & Prod Dev only Enterprise
🍃 Recomandare

Pentru majoritatea proiectelor, Kafka UI (Provectus) este alegerea ideală - oferă un echilibru perfect între funcționalități și ușurință în utilizare, fiind potrivit atât pentru development cât și pentru producție.