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.
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?
Arhitectura de Bază
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. |
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
- Producer trimite un mesaj către un Topic
- Kafka determină Partition-ul (bazat pe key sau round-robin)
- Mesajul primește un Offset unic și este stocat
- 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:
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"
# 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:
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
Accesează start.spring.io sau folosește Spring Initializr din IDE-ul tău preferat.
Adaugă: Spring for Apache Kafka, Spring Web, Lombok (opțional)
<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ă
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ă)
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();
}
}
În producție, setează replicas la minimum 3 pentru toleranță la erori și acks=all pentru durabilitate maximă.
Implementarea Producer-ului
Model de Date
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
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
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
) {}
}
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
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
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
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ă:
@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
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
@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
}
}
Adaugă în application.yml:
spring.kafka.producer.transaction-id-prefix: tx-
Monitoring cu Actuator
management:
endpoints:
web:
exposure:
include: health,metrics,kafka
health:
kafka:
enabled: true
Best Practices
Checklist pentru Producție
- Replication Factor ≥ 3 pentru toleranță la erori
- acks=all pentru durabilitate maximă
- min.insync.replicas=2 pentru a preveni pierderea datelor
- enable.idempotence=true pentru exactly-once semantics
- Configurează retention policy adecvată pentru use case
- Implementează health checks și alerting
- Testează failover scenarios înainte de go-live
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
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ă.
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
# 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
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:
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:
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ă
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
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:
<!-- 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>
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:
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
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 |
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.