Kafka — a distributed event streaming platform designed for high-throughput, fault-tolerant, durable event streaming. RabbitMQ — a traditional message broker implementing AMQP, designed for flexible routing and message queuing.
System Design Context: Kafka is optimized for throughput and event replay — it stores messages durably and allows multiple consumers to read the same events. RabbitMQ is optimized for message routing and delivery — it removes messages after consumption.
Kafka — Detailed Characteristics:
1// Kafka producer with headers and partitioning2@Service3public class EventProducer {45 private final KafkaTemplate<String, OrderEvent> kafkaTemplate;67 public void sendEvent(OrderEvent event) {8 ProducerRecord<String, OrderEvent> record =9 new ProducerRecord<>("order-events", event.orderId().toString(), event);10 record.headers().add("source", "order-service".getBytes());11 record.headers().add("version", "1.0".getBytes());1213 kafkaTemplate.send(record);14 }15}1617// Kafka consumer with manual offset management18@Service19public class EventConsumer {2021 @KafkaListener(22 topics = "order-events",23 groupId = "analytics-service",24 containerFactory = "kafkaListenerContainerFactory"25 )26 public void consume(27 @Payload OrderEvent event,28 Acknowledgment ack) {29 analyticsService.process(event);30 ack.acknowledge(); // manual commit31 }32}
RabbitMQ — Detailed Characteristics:
1// RabbitMQ producer2@Service3public class TaskPublisher {45 private final RabbitTemplate rabbitTemplate;67 public void publishTask(Task task) {8 rabbitTemplate.convertAndSend(9 "task-exchange",10 "task.processing",11 task,12 message -> {13 message.getMessageProperties()14 .setHeader("priority", "high");15 message.getMessageProperties()16 .setExpiration("60000"); // TTL: 60s17 return message;18 });19 }20}2122// RabbitMQ consumer with retry23@Service24public class TaskConsumer {2526 @RabbitListener(27 queues = "task-processing-queue",28 containerFactory = "rabbitListenerContainerFactory"29 )30 public void consume(Task task, Channel channel,31 @Header(AmqpHeaders.DELIVERY_TAG) long tag) {32 try {33 taskService.process(task);34 channel.basicAck(tag, false);35 } catch (Exception e) {36 channel.basicNack(tag, false, true); // requeue37 }38 }39}
Performance Comparison:
Common Pitfalls:
When to Use: