Resilience, Event-Driven Architecture + DevOps: ऐसी Apps बनाना जो कभी नीचे न जाएँ
Production-grade applications के तीन स्तंभों में एक beginner-friendly deep dive: resilience patterns (rate limiting, circuit breakers, bulkheads), Kafka और Spring Events के साथ event-driven messaging, और Prometheus, Grafana, ELK, और Zipkin के साथ DevOps observability।
आपकी App को मज़बूत क्यों होना चाहिए
कल्पना कीजिए कि आप एक pizza की दुकान चलाते हैं। एक दिन, oven टूट जाता है। यदि आपके पास केवल एक oven है, तो किसी को pizza नहीं मिलेगा। लेकिन यदि आपके पास एक backup oven है, oven टूटने पर एक योजना है, और यह जानने का एक तरीका है कि oven टूटने वाला है इससे पहले कि वह हो — तो यही resilience है।
Software में, resilience का मतलब है कि आपकी application तब भी काम करती रहती है जब चीज़ें गलत होती हैं — एक database धीमा हो जाता है, एक partner service crash हो जाती है, या एक करोड़ users एक साथ आ जाते हैं। यह guide तीन स्तंभों को कवर करती है जो apps को लगभग अटूट बनाते हैं: Resilience Patterns, Event-Driven Messaging, और DevOps Observability।
अनुभाग 1 — Resilience: अपनी App को Overload और Failure से बचाना
Rate Limiting क्या है?
Rate limiting एक theme park ride की तरह है — एक बार में केवल 100 लोग ride कर सकते हैं, और बाकी line में प्रतीक्षा करते हैं। Line के बिना, हर कोई एक साथ ride पर दौड़ पड़ता है और यह टूट जाती है। Rate limiting नियंत्रित करती है कि एक user या service किसी दिए गए time window में कितनी requests कर सकता है।
Rate Limiting Algorithms
चार मुख्य algorithms हैं। प्रत्येक को उस theme park line को manage करने के एक अलग तरीके के रूप में सोचें।
1. Fixed Window
समय को बराबर chunks में विभाजित करें (मान लीजिए, 1-मिनट windows)। प्रत्येक window एक निश्चित संख्या में requests की अनुमति देता है। जब window reset होता है, counter शून्य पर वापस चला जाता है।
उपमा: एक candy store प्रत्येक बच्चे को प्रति घंटे 5 candies देती है। हर घंटे की शुरुआत में, count reset हो जाता है — भले ही आपने पहले मिनट में सभी 5 खा ली हों।
Timeline:
|--- Window 1 (00:00-00:59) ---|--- Window 2 (01:00-01:59) ---|
Request 1 ✓ (count: 1) Request 6 ✓ (count: 1)
Request 2 ✓ (count: 2) Request 7 ✓ (count: 2)
Request 3 ✓ (count: 3) ...
Request 4 ✓ (count: 4)
Request 5 ✓ (count: 5)
Request 6 ✗ (REJECTED — limit 5 reached)
Problem: "Boundary burst" — 5 requests at 00:59 + 5 at 01:00 = 10 in 2 seconds!
2. Sliding Window
निश्चित chunks के बजाय, window आपके साथ slide होती है। यह हमेशा अभी से पिछले 60 seconds को देखती है। यह boundary burst समस्या को ठीक करती है।
उपमा: हर घंटे ठीक समय पर candy reset करने के बजाय, store हमेशा देखती है कि "इस बच्चे ने पिछले 60 मिनटों में कितनी candies खाईं?" — समय चाहे जो भी हो।
Current time: 01:30
Look back 60 seconds: 00:30 to 01:30
Count all requests in that range.
If count >= limit → REJECT
If count < limit → ALLOW
Result: Smooth, consistent rate limiting with no burst at boundaries.
3. Token Bucket
कल्पना कीजिए एक bucket जिसमें tokens रखे जाते हैं। हर सेकंड, एक नया token उसमें गिरता है। एक request करने के लिए, आप एक token निकालते हैं। यदि bucket खाली है, तो आप प्रतीक्षा करते हैं। Bucket का एक अधिकतम आकार होता है, इसलिए tokens हमेशा के लिए जमा नहीं होते।
उपमा: एक gumball machine हर 10 seconds में एक gumball refill करती है। आप किसी भी समय एक को पकड़ सकते हैं जब एक gumball हो। यदि machine खाली है, तो आप प्रतीक्षा करते हैं। लेकिन यह कभी भी 10 से अधिक gumballs नहीं रखती, भले ही कोई इसे कुछ समय के लिए उपयोग न करे।
Bucket: max 10 tokens, refill 1 token/second
Time 0s: Bucket = 10 tokens
Time 0s: Burst of 10 requests → all succeed, bucket = 0
Time 1s: 1 token refilled → bucket = 1 → 1 request succeeds
Time 2s: 1 token refilled → bucket = 1 → 1 request succeeds
Time 10s: If idle, bucket = 10 again (back to full)
Key benefit: Allows short bursts while maintaining average rate.
4. Leaky Bucket
Requests पानी की तरह ऊपर से आती हैं। Bucket उन्हें नीचे से एक निश्चित दर पर "leak" (process) करती है। यदि पानी बहुत तेज़ी से आता है, तो bucket overflow हो जाती है और अतिरिक्त requests गिर जाती हैं।
उपमा: एक पानी की बोतल पर एक funnel — इससे कोई फर्क नहीं पड़ता कि आप कितनी तेज़ी से पानी डालते हैं, यह एक स्थिर गति से टपकता रहता है। बहुत तेज़ी से डालो और यह overflow हो जाती है।
Leaky Bucket: capacity 10, leak rate 2/second
Incoming: 5 requests arrive at once
Bucket: [■ ■ ■ ■ ■ · · · · ·] (5/10 — all queued)
Output: 2 requests processed per second
Incoming: 8 more requests arrive
Bucket: [■ ■ ■ ■ ■ ■ ■ ■ ■ ■] (10/10 — full!)
Next req: DROPPED (bucket overflow)
Output: Still processing at steady 2/second
Key benefit: Perfectly smooth output rate, no bursts at all.
Algorithm तुलना तालिका
| Algorithm | Burst Handling | Memory | Accuracy | Best For |
|---|---|---|---|---|
| Fixed Window | Boundary burst समस्या | कम | Medium | Simple APIs, कम traffic |
| Sliding Window | कोई burst समस्या नहीं | Medium | High | Precision की आवश्यकता वाले Production APIs |
| Token Bucket | नियंत्रित bursts की अनुमति देता है | कम | High | Burst tolerance की आवश्यकता वाले APIs |
| Leaky Bucket | कोई bursts नहीं — smooth output | कम | High | Steady-rate processing (queues) |
Circuit Breaker Pattern
एक circuit breaker ठीक उसी तरह काम करता है जैसा आपके घर में होता है। यदि बहुत अधिक बिजली बहती है, तो breaker trip हो जाता है और आग को रोकने के लिए power काट देता है। Software में, यदि आप जिस service को call करते हैं वह बार-बार विफल होती रहती है, तो circuit breaker "trip" हो जाता है और इसे call करना बंद कर देता है — ताकि आपका पूरा system एक मरी हुई service की प्रतीक्षा करते हुए crash न हो।
Circuit Breaker States:
┌────────┐ failures >= threshold ┌──────────┐
│ CLOSED │ ─────────────────────────→ │ OPEN │
│(normal)│ │(blocking)│
└────────┘ └──────────┘
↑ │
│ wait timeout expires │
│ ▼
│ ┌──────────────┐
└────── success ───────────────│ HALF-OPEN │
│(testing 1 req)│
└──────────────┘
│
failure → back to OPEN
CLOSED — सब कुछ सामान्य रूप से काम करता है। Requests pass होती हैं। Breaker विफलताओं की गिनती करता है।
OPEN — बहुत अधिक विफलताएँ! Breaker तुरंत सभी requests को block कर देता है (fast fail)। एक मरी हुई service से timeout के लिए 30 seconds प्रतीक्षा करने की आवश्यकता नहीं।
HALF-OPEN — प्रतीक्षा अवधि के बाद, breaker यह test करने के लिए एक request को pass होने देता है कि क्या service ठीक हुई। यदि यह सफल होती है, तो CLOSED पर वापस जाएँ। यदि यह विफल होती है, तो OPEN पर वापस जाएँ।
Circuit Breaker बनाम Rate Limiter
| पहलू | Circuit Breaker | Rate Limiter |
|---|---|---|
| उद्देश्य | विफल downstream services से सुरक्षा | बहुत अधिक incoming requests से सुरक्षा |
| दिशा | Outgoing calls (आप → दूसरी service) | Incoming calls (user → आप) |
| Trigger होता है | विफलता दर (errors, timeouts) | प्रति time window request count |
| प्रतिक्रिया | Fallback के साथ fast fail | HTTP 429 Too Many Requests |
| उपमा | घर का circuit breaker (आग रोकता है) | Theme park ride line (भीड़ नियंत्रित करती है) |
| साथ उपयोग? | हाँ — incoming को rate limit करें, outgoing पर circuit break करें | |
Resilience4j Configuration
Resilience4j Java/Spring Boot में resilience patterns के लिए go-to library है। इसे चार tools के साथ एक toolbox के रूप में सोचें: Circuit Breaker, Retry, Rate Limiter, और Bulkhead।
Circuit Breaker Config
// application.yml
resilience4j:
circuitbreaker:
instances:
paymentService:
slidingWindowSize: 10 # Look at last 10 calls
failureRateThreshold: 50 # Trip if 50% fail
waitDurationInOpenState: 10s # Wait 10s before testing
permittedNumberOfCallsInHalfOpenState: 3 # Test with 3 calls
slidingWindowType: COUNT_BASED
// PaymentClient.java
@Service
public class PaymentClient {
@CircuitBreaker(name = "paymentService", fallbackMethod = "paymentFallback")
public PaymentResponse processPayment(PaymentRequest request) {
return restTemplate.postForObject(
"http://payment-service/api/payments", request, PaymentResponse.class
);
}
// Fallback: runs when circuit is OPEN or call fails
private PaymentResponse paymentFallback(PaymentRequest request, Throwable ex) {
log.warn("Payment service down, queuing for retry: {}", ex.getMessage());
return PaymentResponse.builder()
.status("QUEUED")
.message("Payment will be processed shortly")
.build();
}
}
Retry Config
// application.yml
resilience4j:
retry:
instances:
inventoryService:
maxAttempts: 3 # Try 3 times total
waitDuration: 500ms # Wait 500ms between retries
retryExceptions:
- java.io.IOException # Retry on network errors
- java.util.concurrent.TimeoutException
ignoreExceptions:
- com.example.BusinessException # Don't retry business errors
// InventoryClient.java
@Service
public class InventoryClient {
@Retry(name = "inventoryService", fallbackMethod = "inventoryFallback")
public InventoryResponse checkStock(String productId) {
return restTemplate.getForObject(
"http://inventory-service/api/stock/" + productId,
InventoryResponse.class
);
}
private InventoryResponse inventoryFallback(String productId, Throwable ex) {
log.warn("Inventory check failed after retries for product {}: {}",
productId, ex.getMessage());
return InventoryResponse.builder()
.productId(productId)
.available(false)
.message("Stock status temporarily unavailable")
.build();
}
}
Rate Limiter Config
// application.yml
resilience4j:
ratelimiter:
instances:
orderApi:
limitForPeriod: 100 # 100 requests allowed
limitRefreshPeriod: 1s # Per 1 second
timeoutDuration: 500ms # Wait up to 500ms for a permit
// OrderController.java
@RestController
@RequestMapping("/api/orders")
public class OrderController {
@RateLimiter(name = "orderApi")
@PostMapping
public ResponseEntity<OrderResponse> createOrder(@RequestBody OrderRequest request) {
return ResponseEntity.ok(orderService.createOrder(request));
}
// If limit exceeded → RequestNotPermitted exception → HTTP 429
}
Bulkhead Pattern
एक bulkhead एक जहाज़ में जलरोधक compartments की तरह है। यदि एक compartment में पानी भर जाता है, तो बाकी सूखे रहते हैं और जहाज़ तैरता रहता है। Software में, एक bulkhead सीमित करता है कि एक operation कितने threads (या concurrent calls) का उपयोग कर सकता है — ताकि एक धीमी service आपके सभी threads को खा न सके और बाकी सब कुछ भूखा न मार सके।
// application.yml
resilience4j:
bulkhead:
instances:
reportService:
maxConcurrentCalls: 5 # Only 5 threads at a time
maxWaitDuration: 100ms # Wait max 100ms for a slot
// ReportController.java
@RestController
public class ReportController {
@Bulkhead(name = "reportService", fallbackMethod = "reportFallback")
@GetMapping("/api/reports/{id}")
public ReportResponse getReport(@PathVariable String id) {
// Even if report generation is slow, only 5 threads used at most
return reportService.generate(id);
}
private ReportResponse reportFallback(String id, Throwable ex) {
return ReportResponse.builder()
.status("BUSY")
.message("Report service is at capacity, please try again shortly")
.build();
}
}
चारों को एक साथ जोड़ना
// You can stack annotations — they execute in this order:
// Retry → CircuitBreaker → RateLimiter → Bulkhead → Your method
@Retry(name = "orderService")
@CircuitBreaker(name = "orderService", fallbackMethod = "fallback")
@RateLimiter(name = "orderService")
@Bulkhead(name = "orderService")
public OrderResponse placeOrder(OrderRequest request) {
return restTemplate.postForObject(
"http://order-service/api/orders", request, OrderResponse.class
);
}
अनुभाग 2 — Kafka के साथ Event-Driven Architecture और Messaging
Event-Driven Architecture क्या है?
पारंपरिक apps में, services एक दूसरे को सीधे call करती हैं: "अरे order-service, मुझे order details चाहिए!" यदि order-service down है, तो caller अटक जाता है।
Event-driven architecture में, services events (messages) publish करके संवाद करती हैं। कोई किसी को सीधे call नहीं करता। इसके बजाय, वे एक साझा mailbox में एक message डालते हैं और जो भी interested है वह इसे उठा लेता है।
उपमा: अपने दोस्त को फोन पर call करने के बजाय (synchronous — उन्हें उठाना होगा), आप post office के माध्यम से एक पत्र भेजते हैं (asynchronous — वे इसे तब पढ़ते हैं जब वे कर सकते हैं)। यदि आपका दोस्त छुट्टी पर है, तो पत्र उनके mailbox में प्रतीक्षा करता है।
Kafka मूल अवधारणाएँ
Apache Kafka event-driven systems के लिए सबसे लोकप्रिय message broker है। Kafka को एक विशाल post office के रूप में सोचें। यहाँ इसके प्रमुख भाग हैं:
1. Topic
एक topic एक named mailbox है। आप विभिन्न प्रकार के events के लिए topics बनाते हैं: order-events, payment-events, inventory-alerts। Producers messages को IN रखते हैं, consumers messages को OUT निकालते हैं।
2. Partition
प्रत्येक topic partitions में विभाजित होता है (0, 1, 2, ... संख्याओं वाले)। Partitions parallel processing की अनुमति देते हैं — कई consumers एक ही समय में विभिन्न partitions से पढ़ सकते हैं।
Topic: order-events (3 partitions)
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Partition 0 │ │ Partition 1 │ │ Partition 2 │
│ [msg0][msg3] │ │ [msg1][msg4] │ │ [msg2][msg5] │
│ [msg6]... │ │ [msg7]... │ │ [msg8]... │
└─────────────┘ └─────────────┘ └─────────────┘
Messages with the same KEY always go to the same partition.
→ All events for order-123 go to the same partition (ordered!).
3. Consumer Group
एक consumer group consumers की एक team है जो एक साथ काम करती है। Kafka प्रत्येक partition को group में ठीक एक consumer को assign करता है। इसका मतलब है कि प्रत्येक message group में केवल एक बार process होता है।
Consumer Group: "order-processor"
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Partition 0 │ │ Partition 1 │ │ Partition 2 │
│ ↓ │ │ ↓ │ │ ↓ │
│ Consumer A │ │ Consumer B │ │ Consumer C │
└─────────────┘ └─────────────┘ └─────────────┘
If Consumer B crashes → Kafka re-assigns Partition 1 to A or C.
If you add Consumer D → Kafka rebalances (one consumer may go idle
if there are more consumers than partitions).
4. Offset
एक offset एक bookmark है। यह Kafka को बताता है: "मैंने इस partition में message #47 तक पढ़ा है।" यदि एक consumer crash होकर restart होता है, तो यह अंतिम committed offset से उठा लेता है — कोई messages खोए नहीं, कोई duplicates नहीं।
Partition 0: [msg0] [msg1] [msg2] [msg3] [msg4] [msg5]
↑
committed offset = 3
Consumer will read msg3 next after restart.
5. Key
Message key यह निर्धारित करता है कि message किस partition में जाता है। समान key वाले messages हमेशा एक ही partition में पहुँचते हैं। यह संबंधित events के लिए ordering की गारंटी देता है।
Key = "order-123" → hash("order-123") % 3 = Partition 1
Key = "order-456" → hash("order-456") % 3 = Partition 0
Key = null → Round-robin across partitions
KafkaTemplate के साथ Kafka Producer
// KafkaProducerConfig.java
@Configuration
public class KafkaProducerConfig {
@Bean
public ProducerFactory<String, String> producerFactory() {
Map<String, Object> config = new HashMap<>();
config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
config.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
config.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
config.put(ProducerConfig.ACKS_CONFIG, "all"); // Wait for all replicas
return new DefaultKafkaProducerFactory<>(config);
}
@Bean
public KafkaTemplate<String, String> kafkaTemplate() {
return new KafkaTemplate<>(producerFactory());
}
}
// OrderEventPublisher.java
@Service
@RequiredArgsConstructor
@Slf4j
public class OrderEventPublisher {
private final KafkaTemplate<String, String> kafkaTemplate;
private final ObjectMapper objectMapper;
public void publishOrderCreated(Order order) {
try {
OrderEvent event = OrderEvent.builder()
.eventType("ORDER_CREATED")
.orderId(order.getId())
.customerId(order.getCustomerId())
.totalAmount(order.getTotalAmount())
.timestamp(Instant.now())
.build();
String payload = objectMapper.writeValueAsString(event);
// Key = orderId → all events for same order go to same partition
kafkaTemplate.send("order-events", order.getId(), payload)
.whenComplete((result, ex) -> {
if (ex == null) {
log.info("Published ORDER_CREATED for order {}. Partition: {}, Offset: {}",
order.getId(),
result.getRecordMetadata().partition(),
result.getRecordMetadata().offset());
} else {
log.error("Failed to publish event for order {}: {}",
order.getId(), ex.getMessage());
}
});
} catch (JsonProcessingException e) {
log.error("Failed to serialize order event: {}", e.getMessage());
}
}
}
@KafkaListener के साथ Kafka Consumer
// KafkaConsumerConfig.java
@Configuration
public class KafkaConsumerConfig {
@Bean
public ConsumerFactory<String, String> consumerFactory() {
Map<String, Object> config = new HashMap<>();
config.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
config.put(ConsumerConfig.GROUP_ID_CONFIG, "payment-processor");
config.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
config.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
config.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); // Manual commit
return new DefaultKafkaConsumerFactory<>(config);
}
@Bean
public ConcurrentKafkaListenerContainerFactory<String, String>
kafkaListenerContainerFactory() {
ConcurrentKafkaListenerContainerFactory<String, String> factory =
new ConcurrentKafkaListenerContainerFactory<>();
factory.setConsumerFactory(consumerFactory());
factory.getContainerProperties()
.setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE);
return factory;
}
}
// PaymentEventListener.java
@Service
@Slf4j
public class PaymentEventListener {
@KafkaListener(
topics = "order-events",
groupId = "payment-processor",
concurrency = "3" // 3 threads → can consume from 3 partitions in parallel
)
public void handleOrderEvent(
@Payload String message,
@Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
@Header(KafkaHeaders.OFFSET) long offset,
Acknowledgment acknowledgment) {
try {
OrderEvent event = objectMapper.readValue(message, OrderEvent.class);
log.info("Received {} from partition {} at offset {}",
event.getEventType(), partition, offset);
if ("ORDER_CREATED".equals(event.getEventType())) {
paymentService.initiatePayment(event.getOrderId(), event.getTotalAmount());
}
// Manually acknowledge — Kafka saves the offset
acknowledgment.acknowledge();
} catch (Exception e) {
log.error("Failed to process message at partition {} offset {}: {}",
partition, offset, e.getMessage());
// Don't acknowledge — message will be redelivered
}
}
}
Event-Driven Patterns
1. Event Sourcing
केवल वर्तमान state store करने के बजाय (जैसे "order total = $50"), आप हर event को store करते हैं जो हुआ: "item added $20", "item added $30", "coupon applied -$5"। वर्तमान state सभी events को replay करके फिर से बनाया जाता है।
उपमा: अपने bank balance ($500) को देखने के बजाय, आप हर transaction रखते हैं: +$1000 salary, -$200 rent, -$300 groceries। आप हमेशा balance को फिर से calculate कर सकते हैं और साथ ही "मैंने पिछले महीने groceries पर कितना खर्च किया?" जैसे सवालों का जवाब दे सकते हैं।
// Events stored in order:
OrderCreated { orderId: "123", customerId: "456" }
ItemAdded { orderId: "123", product: "Widget", price: 20.00 }
ItemAdded { orderId: "123", product: "Gadget", price: 30.00 }
CouponApplied { orderId: "123", discount: 5.00 }
OrderConfirmed { orderId: "123", total: 45.00 }
// Replay events → current state:
Order { id: "123", items: [Widget, Gadget], total: 45.00, status: CONFIRMED }
2. CQRS (Command Query Responsibility Segregation)
पढ़ने और लिखने के लिए अलग models का उपयोग करें। Write side commands ("create order", "update stock") को संभालता है। Read side queries ("show me all orders this week") के लिए optimized है।
उपमा: एक restaurant में एक kitchen (write side — खाना बनाता है) और एक menu (read side — दिखाता है कि क्या उपलब्ध है) होता है। Kitchen और menu अलग चीजें हैं, विभिन्न उद्देश्यों के लिए optimized।
Write Side (Commands) Read Side (Queries)
┌─────────────────┐ ┌─────────────────┐
│ OrderCommand │ │ OrderQuery │
│ Handler │───publish event──→│ Handler │
│ │ │ │
│ Normalized DB │ │ Denormalized DB │
│ (3rd normal form)│ │ (flat, fast) │
└─────────────────┘ └─────────────────┘
Benefits:
- Write DB can be PostgreSQL (strong consistency)
- Read DB can be Elasticsearch (fast search)
- Scale reads and writes independently
3. Saga Pattern
एक saga एक ऐसे transaction का management करती है जो कई services में फैली होती है। एक बड़े transaction (जो services के बीच असंभव है) के बजाय, यह local transactions की एक श्रृंखला चलाती है। यदि एक step विफल होता है, तो यह पिछले steps को undo करने के लिए compensating transactions चलाती है।
उपमा: एक छुट्टी book करना: flight reserve करें → hotel reserve करें → car reserve करें। यदि car rental विफल होती है, तो आप hotel रद्द करते हैं, फिर flight रद्द करते हैं — उल्टे क्रम में।
Happy Path:
order-service: Create Order (PENDING)
↓ event: OrderCreated
payment-service: Charge Payment
↓ event: PaymentCompleted
inventory-service: Reserve Stock
↓ event: StockReserved
order-service: Confirm Order (CONFIRMED) ✓
Failure Path (stock unavailable):
order-service: Create Order (PENDING)
↓ event: OrderCreated
payment-service: Charge Payment
↓ event: PaymentCompleted
inventory-service: Reserve Stock → FAILS!
↓ event: StockReservationFailed
payment-service: REFUND Payment (compensating transaction)
↓ event: PaymentRefunded
order-service: Cancel Order (CANCELLED) ✗
4. Outbox Pattern
Outbox pattern एक खतरनाक समस्या को हल करता है: क्या होगा यदि आपकी service database में save करती है लेकिन Kafka event भेजने से पहले crash हो जाती है? Database कहता है "order created" लेकिन कोई event publish नहीं हुआ — system असंगत है।
समाधान: अपने business data के समान database transaction में एक outbox table में event लिखें। एक अलग process outbox table पढ़ता है और Kafka को events publish करता है।
// Step 1: Save order + outbox event in ONE transaction
@Transactional
public Order createOrder(OrderRequest request) {
Order order = orderRepository.save(new Order(request));
// Save event to outbox table (same transaction!)
outboxRepository.save(OutboxEvent.builder()
.aggregateId(order.getId())
.eventType("ORDER_CREATED")
.payload(objectMapper.writeValueAsString(order))
.status("PENDING")
.build());
return order; // Both saved atomically — or both fail
}
// Step 2: Background poller reads outbox and publishes to Kafka
@Scheduled(fixedDelay = 1000)
public void publishOutboxEvents() {
List<OutboxEvent> pending = outboxRepository.findByStatus("PENDING");
for (OutboxEvent event : pending) {
kafkaTemplate.send("order-events", event.getAggregateId(), event.getPayload());
event.setStatus("PUBLISHED");
outboxRepository.save(event);
}
}
Spring Events — In-Process Messaging
हर चीज़ को Kafka की आवश्यकता नहीं है। एक ही application के भीतर events के लिए, Spring के पास built-in event publishing है। यह सरल है और एक service के अंदर code को decouple करने के लिए एकदम सही है।
// 1. Define an event
public record OrderCompletedEvent(String orderId, BigDecimal total, String customerEmail) {}
// 2. Publish the event
@Service
@RequiredArgsConstructor
public class OrderService {
private final ApplicationEventPublisher eventPublisher;
@Transactional
public Order completeOrder(String orderId) {
Order order = orderRepository.findById(orderId).orElseThrow();
order.setStatus("COMPLETED");
orderRepository.save(order);
// Publish event — listeners will handle side effects
eventPublisher.publishEvent(new OrderCompletedEvent(
order.getId(), order.getTotal(), order.getCustomerEmail()
));
return order;
}
}
// 3a. @EventListener — runs immediately (same thread, same transaction)
@Component
@Slf4j
public class InventoryListener {
@EventListener
public void onOrderCompleted(OrderCompletedEvent event) {
log.info("Reducing stock for order {}", event.orderId());
inventoryService.reduceStock(event.orderId());
}
}
// 3b. @TransactionalEventListener — runs AFTER transaction commits
// Safer for side effects like sending emails
@Component
@Slf4j
public class NotificationListener {
@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT)
public void onOrderCompleted(OrderCompletedEvent event) {
log.info("Sending confirmation email for order {}", event.orderId());
emailService.sendConfirmation(event.customerEmail(), event.orderId());
}
}
// Why AFTER_COMMIT? If the transaction rolls back, you don't want
// to have already sent a "your order is confirmed!" email.
अनुभाग 3 — DevOps + Observability: अपनी Running App के अंदर देखना
Spring Boot Actuator — Health और Monitoring Endpoints
Actuator आपकी app को built-in health check endpoints देता है। यह आपकी गाड़ी में dashboard की तरह है — speed, fuel level, engine temperature — लेकिन आपकी application के लिए।
// application.yml — Actuator configuration
management:
endpoints:
web:
exposure:
include: health, info, prometheus, loggers, metrics
endpoint:
health:
show-details: always # Show DB, Redis, Kafka health
probes:
enabled: true # Enable Kubernetes probes
health:
livenessState:
enabled: true
readinessState:
enabled: true
Key Actuator Endpoints
| Endpoint | URL | उद्देश्य |
|---|---|---|
| Health | /actuator/health | समग्र app health (UP/DOWN) + dependencies |
| Liveness | /actuator/health/liveness | क्या app जीवित है? (Kubernetes DOWN होने पर restart करता है) |
| Readiness | /actuator/health/readiness | क्या app traffic संभाल सकती है? (K8s DOWN होने पर requests भेजना बंद करता है) |
| Prometheus | /actuator/prometheus | Prometheus format में metrics (हर 15s पर scrape) |
| Loggers | /actuator/loggers | Runtime पर log levels देखें/बदलें (कोई restart नहीं!) |
// Example: Change log level at runtime (no redeploy!)
// POST /actuator/loggers/com.example.orderservice
// Body: { "configuredLevel": "DEBUG" }
// Example health response:
{
"status": "UP",
"components": {
"db": { "status": "UP", "details": { "database": "PostgreSQL" } },
"redis": { "status": "UP", "details": { "version": "7.2.4" } },
"kafka": { "status": "UP" },
"diskSpace": { "status": "UP", "details": { "free": "42GB" } }
}
}
Deployment Strategies
आप users को नोटिस किए बिना अपनी app को कैसे update करते हैं? चार सामान्य strategies हैं। इसे एक restaurant को renovate करने की तरह सोचें — क्या आप एक हफ्ते के लिए बंद करते हैं, या एक बार में एक कमरे को remodel करते हुए serving जारी रखते हैं?
1. Rolling Update
एक बार में एक instance को बदलें। पुराने और नए versions कुछ समय के लिए साथ-साथ चलते हैं।
Time 0: [v1] [v1] [v1] [v1] ← 4 instances running v1
Time 1: [v2] [v1] [v1] [v1] ← Replace first instance
Time 2: [v2] [v2] [v1] [v1] ← Replace second
Time 3: [v2] [v2] [v2] [v1] ← Replace third
Time 4: [v2] [v2] [v2] [v2] ← All running v2 ✓
2. Blue-Green Deployment
दो समान environments चलाएँ: Blue (current) और Green (new)। Green पर deploy करें, इसका test करें, फिर सभी traffic को Blue से Green पर तुरंत switch करें।
Step 1: Blue (v1) ← ALL TRAFFIC Green (v2) testing...
Step 2: Blue (v1) Green (v2) ← ALL TRAFFIC (switch!)
Step 3: Blue (standby for rollback) Green (v2) ← serving users
Rollback? Switch traffic back to Blue in seconds!
3. Canary Deployment
Traffic का एक छोटा प्रतिशत (मान लीजिए 5%) नए version पर भेजें। Errors के लिए monitor करें। यदि सब कुछ अच्छा दिखता है, तो धीरे-धीरे 100% तक बढ़ाएँ।
Step 1: v1 ← 100% traffic v2 ← 0%
Step 2: v1 ← 95% traffic v2 ← 5% (canary — watching metrics)
Step 3: v1 ← 70% traffic v2 ← 30% (metrics look good!)
Step 4: v1 ← 0% v2 ← 100% (full rollout ✓)
If errors spike at 5%: kill the canary, 100% back to v1. Only 5% of users affected.
4. Recreate
सब कुछ बंद करें, नया version deploy करें, सब कुछ शुरू करें। सरल लेकिन downtime का कारण बनता है।
Step 1: [v1] [v1] [v1] ← running
Step 2: [ ] [ ] [ ] ← all stopped (DOWNTIME!)
Step 3: [v2] [v2] [v2] ← all started with new version
Deployment Strategy तुलना
| Strategy | Downtime | Risk | Rollback Speed | Resource Cost | Best For |
|---|---|---|---|---|---|
| Rolling | Zero | Medium | धीमा (एक-एक करके roll back) | कम (समान infra) | अधिकांश applications (default) |
| Blue-Green | Zero | कम | तत्काल (वापस switch करें) | High (double infra) | तत्काल rollback की आवश्यकता वाले critical apps |
| Canary | Zero | बहुत कम | तेज़ (canary को kill करें) | Low-Medium | High-traffic apps, नए जोखिम भरे features |
| Recreate | हाँ | High | धीमा (पूर्ण redeploy) | कम | Dev/staging, या downtime सहन कर सकने वाले apps |
Observability Stack — Metrics, Logs, Traces, Alerts
Observability तीन सवालों के जवाब देती है: क्या हुआ? (logs), अभी क्या हो रहा है? (metrics), और यह क्यों हुआ? (traces)। इसे एक doctor की तरह सोचें: लक्षण (metrics), medical history (logs), और एक MRI scan (traces)।
Metrics: Micrometer → Prometheus → Grafana
How it works:
┌──────────────┐ scrape /metrics ┌────────────┐ dashboards ┌─────────┐
│ Spring Boot │ ────────every 15s───→ │ Prometheus │ ──────────────→ │ Grafana │
│ (Micrometer)│ │ (time-series│ │ (charts)│
└──────────────┘ │ database) │ └─────────┘
└────────────┘
// Custom metrics in your code
@Service
@RequiredArgsConstructor
public class OrderMetrics {
private final MeterRegistry meterRegistry;
public void recordOrderPlaced(String orderType) {
// Counter — goes up only (total orders placed)
meterRegistry.counter("orders.placed", "type", orderType).increment();
}
public void recordOrderProcessingTime(long millis) {
// Timer — tracks duration distribution
meterRegistry.timer("orders.processing.time")
.record(millis, TimeUnit.MILLISECONDS);
}
public void recordActiveOrders(int count) {
// Gauge — goes up and down (current value)
meterRegistry.gauge("orders.active.count", count);
}
}
Logging: SLF4J → JSON → ELK Stack
Structured JSON logs machine-readable हैं। ELK stack (Elasticsearch, Logstash, Kibana) उन्हें collect, index, और visualize करती है।
How it works:
┌────────────┐ JSON logs ┌──────────┐ index ┌───────────────┐ search ┌────────┐
│ Spring Boot│ ─────────────→│ Logstash │ ────────→ │ Elasticsearch │ ────────→│ Kibana │
│ (SLF4J) │ (structured) │ (parse) │ │ (store) │ │(search)│
└────────────┘ └──────────┘ └───────────────┘ └────────┘
// logback-spring.xml — JSON structured logging
<configuration>
<appender name="JSON" class="ch.qos.logback.core.ConsoleAppender">
<encoder class="net.logstash.logback.encoder.LogstashEncoder">
<includeMdcKeyName>traceId</includeMdcKeyName>
<includeMdcKeyName>spanId</includeMdcKeyName>
<includeMdcKeyName>userId</includeMdcKeyName>
</encoder>
</appender>
<root level="INFO">
<appender-ref ref="JSON" />
</root>
</configuration>
// Output — one JSON object per log line:
{
"timestamp": "2026-04-08T10:30:00.123Z",
"level": "INFO",
"logger": "com.example.OrderService",
"message": "Order placed successfully",
"traceId": "abc123def456",
"spanId": "789ghi",
"userId": "user-42",
"orderId": "order-789",
"amount": 49.99
}
Tracing: Micrometer Tracing → Zipkin
एक trace एक single request का अनुसरण करता है जब यह कई services के माध्यम से यात्रा करती है। प्रत्येक service एक "span" जोड़ती है (यात्रा का अपना भाग)। आप ठीक-ठीक देख सकते हैं कि समय कहाँ खर्च हुआ।
How it works:
┌──────────┐ ┌──────────┐ ┌──────────┐
│ API │────→│ Order │────→│ Payment │
│ Gateway │ │ Service │ │ Service │
│ span: 2ms│ │ span:45ms│ │ span:30ms│
└──────────┘ └──────────┘ └──────────┘
│ │ │
└────────────────┴────────────────┘
│
▼
┌────────────┐
│ Zipkin │
│ (trace UI) │
└────────────┘
Trace ID: abc-123 (same across all services)
Total time: 77ms
Bottleneck: Payment Service (30ms) — investigate!
// application.yml — Tracing configuration
management:
tracing:
sampling:
probability: 1.0 # Sample 100% in dev, 10% (0.1) in prod
spring:
application:
name: order-service
// pom.xml dependencies:
// micrometer-tracing-bridge-brave
// zipkin-reporter-brave
Alerting: Prometheus AlertManager → Slack/Email
How it works:
┌────────────┐ rule violated ┌──────────────┐ notification ┌───────────┐
│ Prometheus │ ──────────────→ │ AlertManager │ ─────────────→ │ Slack / │
│ (rules) │ │ (routing) │ │ Email │
└────────────┘ └──────────────┘ └───────────┘
# alert-rules.yml — Prometheus alerting rules
groups:
- name: application-alerts
rules:
- alert: HighErrorRate
expr: rate(http_server_requests_seconds_count{status=~"5.."}[5m]) > 10
for: 2m
labels:
severity: critical
annotations:
summary: "High error rate on {{ $labels.instance }}"
description: "More than 10 errors/sec for 2 minutes"
- alert: ServiceDown
expr: up == 0
for: 1m
labels:
severity: critical
annotations:
summary: "{{ $labels.job }} is DOWN"
- alert: HighResponseTime
expr: histogram_quantile(0.95, rate(http_server_requests_seconds_bucket[5m])) > 2
for: 5m
labels:
severity: warning
annotations:
summary: "95th percentile response time > 2 seconds"
# alertmanager.yml — Route alerts to Slack
global:
resolve_timeout: 5m
route:
group_by: [alertname, severity]
group_wait: 10s
group_interval: 5m
receiver: slack-notifications
receivers:
- name: slack-notifications
slack_configs:
- api_url: "https://hooks.slack.com/services/YOUR/WEBHOOK/URL"
channel: "#alerts"
title: '{{ .GroupLabels.alertname }}'
text: '{{ .CommonAnnotations.description }}'
send_resolved: true
पूर्ण Observability Architecture
┌─────────────────────────────────────────────────────────────────┐
│ YOUR APPLICATION │
│ │
│ Micrometer ──→ /actuator/prometheus ──→ Prometheus ──→ Grafana │
│ (metrics) (store) (charts) │
│ │
│ SLF4J ──→ JSON logs ──→ Logstash ──→ Elasticsearch ──→ Kibana │
│ (logging) (parse) (index) (search) │
│ │
│ Micrometer ──→ Trace spans ──→ Zipkin │
│ Tracing (propagated) (trace UI) │
│ │
│ Prometheus ──→ Alert rules ──→ AlertManager ──→ Slack/Email │
│ (thresholds) (routing) (notification) │
└─────────────────────────────────────────────────────────────────┘
अक्सर पूछे जाने वाले प्रश्न
1. मुझे Circuit Breaker बनाम Rate Limiter का उपयोग कब करना चाहिए?
अपनी खुद की service को बहुत अधिक incoming requests से अभिभूत होने से बचाने के लिए Rate Limiter का उपयोग करें (एक club के bouncer की तरह)। अपनी service को एक टूटी हुई downstream service को call करने में समय बर्बाद करने से बचाने के लिए Circuit Breaker का उपयोग करें (एक मरे हुए उपकरण को unplug करने की तरह ताकि यह आपके घर के breaker को trip न करे)। Production में, दोनों एक साथ उपयोग करें: incoming को rate limit करें, outgoing पर circuit break करें।
2. Kafka और Spring Events के बीच क्या अंतर है?
Spring Events एक single application के अंदर काम करती हैं — एक ही कमरे में किसी को एक note passing करने की तरह। Kafka विभिन्न servers पर चलने वाली कई applications में काम करती है — post office के माध्यम से एक पत्र भेजने की तरह। In-process decoupling के लिए Spring Events का उपयोग करें (जैसे order completion के बाद email भेजना)। Kafka का उपयोग तब करें जब विभिन्न machines पर विभिन्न services को asynchronous रूप से संवाद करने की आवश्यकता हो, और आपको message durability की आवश्यकता हो (messages crashes से बच जाते हैं)।
3. Outbox Pattern data असंगति को कैसे रोकता है?
Outbox pattern के बिना, आप database में save करते हैं और फिर Kafka पर publish करते हैं — दो अलग operations। यदि आपकी app उनके बीच crash होती है, तो database में data है लेकिन कोई event publish नहीं हुआ। Outbox pattern business data और event दोनों को एक ही database transaction में लिखता है। या तो दोनों save होते हैं या कोई नहीं। एक अलग background process फिर outbox table पढ़ता है और Kafka को publish करता है। भले ही publisher crash हो जाए, event अभी भी database में सुरक्षित है, अगले run पर publish होने की प्रतीक्षा में।
4. एक production application के लिए best deployment strategy क्या है?
Rolling update सबसे अच्छा default विकल्प है — इसमें zero downtime है, समान infrastructure का उपयोग करता है, और Kubernetes में built-in है। Mission-critical apps के लिए जहाँ आपको तत्काल rollback की आवश्यकता है, Blue-Green का उपयोग करें (लेकिन इसमें double infrastructure खर्च होता है)। High-traffic apps के लिए जहाँ आप सुरक्षित रूप से नए features का test करना चाहते हैं, Canary का उपयोग करें (पहले 5% traffic route करें, फिर धीरे-धीरे बढ़ाएँ)। Recreate केवल development environments या कुछ मिनट के downtime को सहन कर सकने वाले apps के लिए है।
5. Metrics, Logs, और Traces एक साथ कैसे काम करते हैं?
एक production issue की जाँच को एक detective होने जैसा सोचें। Metrics वह alarm है जो आपको बताता है कि कुछ गलत है ("error rate 2:30 PM पर बढ़ गई")। Logs आपको विवरण देते हैं ("user-789 के लिए OrderService.java line 42 में NullPointerException")। Traces आपको services के बीच एक request की पूरी यात्रा दिखाते हैं ("इस request ने order-service में 200ms बिताए, फिर payment-service में 3000ms अटकी रही — यही bottleneck है!")। एक साथ, वे जवाब देते हैं: क्या हुआ, कहाँ हुआ, और क्यों हुआ।