రెసిలియన్స్, ఈవెంట్-డ్రివెన్ ఆర్కిటెక్చర్ + DevOps: ఎప్పటికీ డౌన్ కాని యాప్లను నిర్మించడం
రెసిలియన్స్ patterns (rate limiting, circuit breakers, bulkheads), Kafka మరియు Spring Events తో event-driven messaging, మరియు Prometheus, Grafana, ELK, మరియు Zipkin తో DevOps observability — production-grade applications యొక్క మూడు స్తంభాలలో beginner-friendly deep dive.
మీ యాప్ ఎందుకు గట్టిగా ఉండాలి
మీరు ఒక పిజ్జా షాపు నడుపుతున్నారని ఊహించుకోండి. ఒక రోజు, ఓవెన్ పాడైపోయింది. మీ దగ్గర ఒకే ఓవెన్ ఉంటే, ఎవరికీ పిజ్జా రాదు. కానీ మీ దగ్గర బ్యాకప్ ఓవెన్, ఓవెన్లు పాడైనప్పుడు ప్లాన్, మరియు అది పాడవబోతోందని ముందే తెలుసుకునే మార్గం ఉంటే — అదే resilience.
సాఫ్ట్వేర్లో, resilience అంటే మీ అప్లికేషన్ విషయాలు తప్పుగా జరిగినా పని చేస్తూనే ఉంటుంది — database నెమ్మదిగా అవుతుంది, ఒక సర్వీస్ క్రాష్ అవుతుంది, లేదా ఒకేసారి మిలియన్ యూజర్లు వస్తారు. ఈ గైడ్ యాప్లను దాదాపు బ్రేక్ కాకుండా చేసే మూడు స్తంభాలను కవర్ చేస్తుంది: Resilience Patterns, Event-Driven Messaging, మరియు DevOps Observability.
సెక్షన్ 1 — Resilience: మీ యాప్ను ఓవర్లోడ్ మరియు వైఫల్యం నుండి రక్షించడం
Rate Limiting అంటే ఏమిటి?
Rate limiting అనేది థీమ్ పార్క్ రైడ్ లాంటిది — ఒక సమయంలో 100 మంది మాత్రమే రైడ్ చేయగలరు, మిగతా వారు లైన్లో వేచి ఉంటారు. లైన్ లేకుంటే, అందరూ ఒకేసారి రైడ్ మీదకు పరుగెత్తుతారు మరియు అది విరిగిపోతుంది. Rate limiting ఒక యూజర్ లేదా సర్వీస్ నిర్దిష్ట సమయ వ్యవధిలో ఎన్ని requests చేయగలదో నియంత్రిస్తుంది.
Rate Limiting Algorithms
నాలుగు ప్రధాన algorithms ఉన్నాయి. ప్రతి ఒక్కటి థీమ్ పార్క్ లైన్ను నిర్వహించడానికి వేరే విధానం అని భావించండి.
1. Fixed Window
సమయాన్ని సమాన భాగాలుగా విభజించండి (ఉదాహరణకు, 1 నిమిషం windows). ప్రతి window నిర్ణీత సంఖ్యలో requests ను అనుమతిస్తుంది. Window రీసెట్ అయినప్పుడు, counter తిరిగి zero కి వస్తుంది.
అనాలజీ: ఒక క్యాండీ షాపు ప్రతి పిల్లవాడికి గంటకు 5 క్యాండీలు ఇస్తుంది. ప్రతి గంట ప్రారంభంలో, count రీసెట్ అవుతుంది — మీరు మొదటి నిమిషంలోనే అన్ని 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
స్థిరమైన భాగాలకు బదులుగా, window మీతో జారుతుంది. ఇది ఎల్లప్పుడూ ఇప్పుడు నుండి చివరి 60 సెకన్లను చూస్తుంది. ఇది boundary burst సమస్యను పరిష్కరిస్తుంది.
అనాలజీ: ప్రతి గంటకు సరిగ్గా క్యాండీ రీసెట్ చేయడానికి బదులు, షాపు ఎల్లప్పుడూ "ఈ పిల్లవాడు చివరి 60 నిమిషాల్లో ఎన్ని క్యాండీలు తిన్నాడు?" అని చూస్తుంది — ఏ సమయమైనా సరే.
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
Tokens ను పట్టుకునే ఒక bucket ను ఊహించుకోండి. ప్రతి సెకనుకు, కొత్త token లో పడుతుంది. Request చేయడానికి, మీరు ఒక token ను బయటకు తీస్తారు. Bucket ఖాళీ అయితే, మీరు వేచి ఉంటారు. Bucket కు maximum size ఉంటుంది, కాబట్టి tokens ఎప్పటికీ పేరుకుపోవు.
అనాలజీ: ఒక గమ్బాల్ మెషీన్ ప్రతి 10 సెకన్లకు ఒక గమ్బాల్ రీఫిల్ చేస్తుంది. గమ్బాల్ ఉన్నప్పుడు ఎప్పుడైనా మీరు ఒకటి తీసుకోవచ్చు. మెషీన్ ఖాళీ అయితే, మీరు వేచి ఉంటారు. కానీ ఎవరూ ఉపయోగించకపోయినా ఇది 10 గమ్బాల్స్ కంటే ఎక్కువ పట్టుకోదు.
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 వాటిని అడుగు నుండి స్థిరమైన rate లో "leak" చేస్తుంది (process చేస్తుంది). నీళ్ళు చాలా వేగంగా పోస్తే, bucket ఓవర్ఫ్లో అవుతుంది మరియు అదనపు requests వదిలేయబడతాయి.
అనాలజీ: వాటర్ బాటిల్ మీద ఫన్నెల్ — మీరు ఎంత వేగంగా నీళ్ళు పోసినా, అది ఒకే స్థిరమైన వేగంతో బయటికి చుక్కలు పడుతుంది. చాలా వేగంగా పోస్తే ఓవర్ఫ్లో అవుతుంది.
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 Comparison Table
| Algorithm | Burst Handling | Memory | Accuracy | Best For |
|---|---|---|---|---|
| Fixed Window | Boundary burst issue | Low | Medium | సాధారణ APIs, తక్కువ ట్రాఫిక్ |
| Sliding Window | Burst issues లేవు | Medium | High | ఖచ్చితత్వం అవసరమైన Production APIs |
| Token Bucket | నియంత్రిత bursts ను అనుమతిస్తుంది | Low | High | Burst tolerance కావాల్సిన APIs |
| Leaky Bucket | Bursts లేవు — smooth output | Low | High | Steady-rate processing (queues) |
Circuit Breaker Pattern
Circuit breaker సరిగ్గా మీ ఇంట్లో ఉన్నట్లుగానే పని చేస్తుంది. చాలా ఎక్కువ విద్యుత్తు ప్రవహిస్తే, breaker trip అవుతుంది మరియు అగ్ని ప్రమాదాన్ని నివారించడానికి విద్యుత్తును కట్ చేస్తుంది. సాఫ్ట్వేర్లో, మీరు call చేసే సర్వీస్ వరుసగా fail అవుతూ ఉంటే, circuit breaker "trip" అవుతుంది మరియు దానిని call చేయడం ఆపేస్తుంది — తద్వారా మీ మొత్తం system చనిపోయిన service కోసం wait చేస్తూ 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 through అవుతాయి. Breaker failures ను count చేస్తుంది.
OPEN — చాలా ఎక్కువ failures! Breaker అన్ని requests ను వెంటనే block చేస్తుంది (fast fail). చనిపోయిన service నుండి 30 సెకన్ల timeout కోసం wait చేయడం ఇక లేదు.
HALF-OPEN — Wait period తర్వాత, service recover అయిందా అని test చేయడానికి breaker ఒక request ను let through చేస్తుంది. అది succeed అయితే, CLOSED కి తిరిగి వెళ్ళండి. Fail అయితే, OPEN కి తిరిగి వెళ్ళండి.
Circuit Breaker vs Rate Limiter
| అంశం | Circuit Breaker | Rate Limiter |
|---|---|---|
| ఉద్దేశం | Fail అవుతున్న downstream services నుండి రక్షణ | చాలా ఎక్కువ incoming requests నుండి రక్షణ |
| దిశ | Outgoing calls (మీరు → ఇతర service) | Incoming calls (user → మీరు) |
| Trigger | Failure rate (errors, timeouts) | Time window కు Request count |
| Response | Fallback తో Fast fail | HTTP 429 Too Many Requests |
| అనాలజీ | ఇంటి circuit breaker (అగ్ని నివారణ) | థీమ్ పార్క్ రైడ్ లైన్ (crowd నియంత్రణ) |
| కలిపి ఉపయోగించాలా? | అవును — incoming ను rate limit చేయండి, outgoing ను circuit break చేయండి | |
Resilience4j Configuration
Resilience4j అనేది Java/Spring Boot లో resilience patterns కోసం ప్రధాన library. దీన్ని నాలుగు tools ఉన్న toolbox గా భావించండి: Circuit Breaker, Retry, Rate Limiter, మరియు Bulkhead.
Circuit Breaker Config
// application.yml
resilience4j:
circuitbreaker:
instances:
paymentService:
slidingWindowSize: 10 # చివరి 10 calls చూడండి
failureRateThreshold: 50 # 50% fail అయితే trip అవుతుంది
waitDurationInOpenState: 10s # Test చేయడానికి ముందు 10s wait
permittedNumberOfCallsInHalfOpenState: 3 # 3 calls తో test
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: circuit OPEN అయినప్పుడు లేదా call fail అయినప్పుడు run అవుతుంది
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 # మొత్తం 3 సార్లు ప్రయత్నించండి
waitDuration: 500ms # Retries మధ్య 500ms wait
retryExceptions:
- java.io.IOException # Network errors మీద retry
- java.util.concurrent.TimeoutException
ignoreExceptions:
- com.example.BusinessException # Business errors మీద retry చేయకండి
// 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 అనుమతించబడతాయి
limitRefreshPeriod: 1s # ప్రతి 1 సెకనుకు
timeoutDuration: 500ms # Permit కోసం 500ms వరకు wait
// 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));
}
// Limit మించితే → RequestNotPermitted exception → HTTP 429
}
Bulkhead Pattern
Bulkhead అనేది ఓడలో waterproof compartments లాంటిది. ఒక compartment నీళ్ళతో నిండితే, మిగతావి dry గా ఉంటాయి మరియు ఓడ తేలుతూ ఉంటుంది. సాఫ్ట్వేర్లో, bulkhead ఒక operation ఎన్ని threads (లేదా concurrent calls) ఉపయోగించగలదో పరిమితం చేస్తుంది — కాబట్టి నెమ్మదిగా ఉండే service మీ అన్ని threads ను తినేసి అన్నింటినీ starve చేయదు.
// application.yml
resilience4j:
bulkhead:
instances:
reportService:
maxConcurrentCalls: 5 # ఒక సమయంలో 5 threads మాత్రమే
maxWaitDuration: 100ms # Slot కోసం max 100ms wait
// ReportController.java
@RestController
public class ReportController {
@Bulkhead(name = "reportService", fallbackMethod = "reportFallback")
@GetMapping("/api/reports/{id}")
public ReportResponse getReport(@PathVariable String id) {
// Report generation నెమ్మదిగా ఉన్నా, గరిష్టంగా 5 threads మాత్రమే ఉపయోగించబడతాయి
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();
}
}
నాలుగింటినీ కలిపి ఉపయోగించడం
// మీరు annotations ను stack చేయవచ్చు — అవి ఈ క్రమంలో execute అవుతాయి:
// 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 — Event-Driven Architecture & Kafka తో Messaging
Event-Driven Architecture అంటే ఏమిటి?
సాంప్రదాయ apps లో, services ఒకదానిని ఒకటి నేరుగా call చేస్తాయి: "హేయ్ order-service, నాకు order details కావాలి!" order-service down అయితే, caller stuck అవుతుంది.
Event-driven architecture లో, services events (messages) publish చేయడం ద్వారా communicate చేస్తాయి. ఎవరూ ఎవరినీ నేరుగా call చేయరు. బదులుగా, వారు shared mailbox లో message drop చేస్తారు మరియు ఎవరు interested ఉన్నారో వారు దానిని pick up చేస్తారు.
అనాలజీ: మీ friend కి phone లో call చేయడానికి బదులు (synchronous — వారు pick up చేయాలి), మీరు post office ద్వారా letter పంపుతారు (asynchronous — వారు చదవగలిగినప్పుడు చదువుతారు). మీ friend vacation లో ఉంటే, letter వారి mailbox లో wait చేస్తుంది.
Kafka Core Concepts
Apache Kafka event-driven systems కోసం అత్యంత ప్రాచుర్యమైన message broker. Kafka ను ఒక పెద్ద post office గా భావించండి. ఇక్కడ దాని ముఖ్యమైన భాగాలు ఉన్నాయి:
1. Topic
Topic అనేది ఒక పేరు ఉన్న mailbox. మీరు వివిధ రకాల events కోసం topics సృష్టిస్తారు: order-events, payment-events, inventory-alerts. Producers messages ను IN పెడతారు, consumers messages ను OUT తీస్తారు.
2. Partition
ప్రతి topic partitions గా విభజించబడుతుంది (0, 1, 2, ... అని number చేయబడతాయి). Partitions parallel processing ను అనుమతిస్తాయి — multiple consumers ఒకే సమయంలో వేర్వేరు partitions నుండి read చేయగలరు.
Topic: order-events (3 partitions)
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Partition 0 │ │ Partition 1 │ │ Partition 2 │
│ [msg0][msg3] │ │ [msg1][msg4] │ │ [msg2][msg5] │
│ [msg6]... │ │ [msg7]... │ │ [msg8]... │
└─────────────┘ └─────────────┘ └─────────────┘
ఒకే KEY ఉన్న Messages ఎల్లప్పుడూ ఒకే partition కి వెళ్తాయి.
→ order-123 కోసం అన్ని events ఒకే partition కి వెళ్తాయి (ordered!).
3. Consumer Group
Consumer group అనేది కలిసి పని చేసే consumers బృందం. Kafka ప్రతి partition ను group లో ఒక్క consumer కి assign చేస్తుంది. అంటే ప్రతి message group కి ఒకసారి మాత్రమే process అవుతుంది.
Consumer Group: "order-processor"
┌─────────────┐ ┌─────────────┐ ┌─────────────┐
│ Partition 0 │ │ Partition 1 │ │ Partition 2 │
│ ↓ │ │ ↓ │ │ ↓ │
│ Consumer A │ │ Consumer B │ │ Consumer C │
└─────────────┘ └─────────────┘ └─────────────┘
Consumer B crash అయితే → Kafka Partition 1 ను A లేదా C కి re-assign చేస్తుంది.
Consumer D add చేస్తే → Kafka rebalance చేస్తుంది (partitions కంటే ఎక్కువ
consumers ఉంటే ఒక consumer idle గా ఉండవచ్చు).
4. Offset
Offset అనేది ఒక bookmark. ఇది Kafka కి చెబుతుంది: "నేను ఈ partition లో message #47 వరకు read చేసాను." Consumer crash అయి restart అయితే, ఇది చివరిగా committed offset నుండి pick up చేస్తుంది — messages కోల్పోవు, duplicates లేవు.
Partition 0: [msg0] [msg1] [msg2] [msg3] [msg4] [msg5]
↑
committed offset = 3
Restart తర్వాత Consumer msg3 ను next read చేస్తుంది.
5. Key
Message key ఏ partition కి message వెళ్తుందో నిర్ణయిస్తుంది. ఒకే key ఉన్న messages ఎల్లప్పుడూ ఒకే partition లో land అవుతాయి. ఇది related events కోసం ordering guarantee ఇస్తుంది.
Key = "order-123" → hash("order-123") % 3 = Partition 1
Key = "order-456" → hash("order-456") % 3 = Partition 0
Key = null → Partitions అంతటా Round-robin
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"); // అన్ని replicas కోసం wait
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 → ఒకే order కి అన్ని events ఒకే 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 → 3 partitions నుండి parallel గా consume చేయగలవు
)
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 offset ను save చేస్తుంది
acknowledgment.acknowledge();
} catch (Exception e) {
log.error("Failed to process message at partition {} offset {}: {}",
partition, offset, e.getMessage());
// Acknowledge చేయకండి — message redeliver అవుతుంది
}
}
}
Event-Driven Patterns
1. Event Sourcing
ప్రస్తుత state ను మాత్రమే store చేయడానికి బదులు ("order total = $50"), జరిగిన ప్రతి event ను store చేయండి: "item added $20", "item added $30", "coupon applied -$5". అన్ని events ను replay చేయడం ద్వారా ప్రస్తుత state rebuild అవుతుంది.
అనాలజీ: మీ bank balance ($500) ను చూడటానికి బదులు, మీరు ప్రతి transaction ను keep చేస్తారు: +$1000 salary, -$200 rent, -$300 groceries. మీరు ఎల్లప్పుడూ balance ను recalculate చేయగలరు మరియు "గత నెలలో 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 }
// Events replay → current state:
Order { id: "123", items: [Widget, Gadget], total: 45.00, status: CONFIRMED }
2. CQRS (Command Query Responsibility Segregation)
Reading మరియు writing కోసం separate models ఉపయోగించండి. Write side commands ను handle చేస్తుంది ("create order", "update stock"). Read side queries కోసం optimize అయి ఉంటుంది ("ఈ వారంలో అన్ని orders చూపించండి").
అనాలజీ: ఒక restaurant లో kitchen (write side — food create చేస్తుంది) మరియు menu (read side — ఏమి available ఉందో చూపిస్తుంది) ఉంటాయి. Kitchen మరియు menu వేర్వేరు విషయాలు, వేర్వేరు purposes కోసం optimize అయి ఉంటాయి.
Write Side (Commands) Read Side (Queries)
┌─────────────────┐ ┌─────────────────┐
│ OrderCommand │ │ OrderQuery │
│ Handler │───publish event──→│ Handler │
│ │ │ │
│ Normalized DB │ │ Denormalized DB │
│ (3rd normal form)│ │ (flat, fast) │
└─────────────────┘ └─────────────────┘
Benefits:
- Write DB PostgreSQL (strong consistency) కావచ్చు
- Read DB Elasticsearch (fast search) కావచ్చు
- Reads మరియు writes ను independently scale చేయవచ్చు
3. Saga Pattern
Saga multiple services అంతటా విస్తరించిన transaction ను manage చేస్తుంది. ఒక పెద్ద transaction (services అంతటా impossible) బదులు, ఇది local transactions chain ను run చేస్తుంది. ఒక step fail అయితే, ఇది మునుపటి steps ను undo చేయడానికి compensating transactions ను run చేస్తుంది.
అనాలజీ: Vacation booking: flight reserve → hotel reserve → car reserve. Car rental fail అయితే, మీరు hotel cancel చేస్తారు, తర్వాత flight cancel చేస్తారు — reverse order లో.
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 inconsistent అయింది.
పరిష్కారం: మీ business data తో SAME database transaction లో event ను outbox table కి write చేయండి. ఒక separate process outbox table ను read చేసి events ను Kafka కి publish చేస్తుంది.
// Step 1: ఒకే transaction లో order + outbox event save చేయండి
@Transactional
public Order createOrder(OrderRequest request) {
Order order = orderRepository.save(new Order(request));
// Outbox table కి event save చేయండి (same transaction!)
outboxRepository.save(OutboxEvent.builder()
.aggregateId(order.getId())
.eventType("ORDER_CREATED")
.payload(objectMapper.writeValueAsString(order))
.status("PENDING")
.build());
return order; // రెండూ atomically save అవుతాయి — లేదా రెండూ fail
}
// Step 2: Background poller outbox read చేసి Kafka కి publish చేస్తుంది
@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 చేయడానికి perfect.
// 1. Event ను define చేయండి
public record OrderCompletedEvent(String orderId, BigDecimal total, String customerEmail) {}
// 2. Event ను publish చేయండి
@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);
// Event publish చేయండి — listeners side effects handle చేస్తాయి
eventPublisher.publishEvent(new OrderCompletedEvent(
order.getId(), order.getTotal(), order.getCustomerEmail()
));
return order;
}
}
// 3a. @EventListener — వెంటనే run అవుతుంది (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 — transaction commit అయిన తర్వాత run అవుతుంది
// Emails పంపడం వంటి side effects కోసం safer
@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());
}
}
// AFTER_COMMIT ఎందుకు? Transaction rollback అయితే, మీరు ఇప్పటికే
// "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 # DB, Redis, Kafka health చూపించండి
probes:
enabled: true # Kubernetes probes enable చేయండి
health:
livenessState:
enabled: true
readinessState:
enabled: true
ముఖ్యమైన Actuator Endpoints
| Endpoint | URL | ఉద్దేశం |
|---|---|---|
| Health | /actuator/health | Overall app health (UP/DOWN) + dependencies |
| Liveness | /actuator/health/liveness | App alive ఉందా? (DOWN అయితే Kubernetes restart చేస్తుంది) |
| Readiness | /actuator/health/readiness | App traffic handle చేయగలదా? (DOWN అయితే K8s requests పంపడం ఆపుతుంది) |
| Prometheus | /actuator/prometheus | Prometheus format లో metrics (ప్రతి 15s scrape అవుతాయి) |
| Loggers | /actuator/loggers | Runtime లో log levels చూడండి/మార్చండి (restart అవసరం లేదు!) |
// ఉదాహరణ: Runtime లో log level మార్చండి (redeploy అవసరం లేదు!)
// POST /actuator/loggers/com.example.orderservice
// Body: { "configuredLevel": "DEBUG" }
// ఉదాహరణ 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 చేయడం లాగా ఆలోచించండి — మీరు ఒక వారం మూసేస్తారా, లేదా ఒక room ను remodel చేస్తూనే serve చేస్తూ ఉంటారా?
1. Rolling Update
Instances ను ఒకదాని తర్వాత ఒకటి replace చేయండి. పాత మరియు కొత్త versions కొద్ది సమయం side by side run అవుతాయి.
Time 0: [v1] [v1] [v1] [v1] ← 4 instances v1 run అవుతున్నాయి
Time 1: [v2] [v1] [v1] [v1] ← మొదటి instance replace
Time 2: [v2] [v2] [v1] [v1] ← రెండవది replace
Time 3: [v2] [v2] [v2] [v1] ← మూడవది replace
Time 4: [v2] [v2] [v2] [v2] ← అన్నీ v2 run అవుతున్నాయి ✓
2. Blue-Green Deployment
రెండు identical environments run చేయండి: Blue (ప్రస్తుతం) మరియు Green (కొత్తది). Green కి deploy చేయండి, test చేయండి, తర్వాత అన్ని traffic ను Blue నుండి Green కి instantly switch చేయండి.
Step 1: Blue (v1) ← ALL TRAFFIC Green (v2) testing...
Step 2: Blue (v1) Green (v2) ← ALL TRAFFIC (switch!)
Step 3: Blue (rollback కోసం standby) Green (v2) ← users serve చేస్తుంది
Rollback? సెకన్లలో traffic తిరిగి Blue కి switch చేయండి!
3. Canary Deployment
Traffic లో చిన్న శాతం (ఉదాహరణకు 5%) కొత్త version కి పంపండి. Errors కోసం monitor చేయండి. అంతా బాగుంటే, క్రమంగా 100% కి పెంచండి.
Step 1: v1 ← 100% traffic v2 ← 0%
Step 2: v1 ← 95% traffic v2 ← 5% (canary — metrics watch చేస్తుంది)
Step 3: v1 ← 70% traffic v2 ← 30% (metrics బాగున్నాయి!)
Step 4: v1 ← 0% v2 ← 100% (full rollout ✓)
5% వద్ద errors spike అయితే: canary ను kill చేయండి, 100% v1 కి తిరిగి. 5% users మాత్రమే affected.
4. Recreate
అన్నీ ఆపండి, కొత్త version deploy చేయండి, అన్నీ start చేయండి. సరళమైనది కానీ downtime కలిగిస్తుంది.
Step 1: [v1] [v1] [v1] ← running
Step 2: [ ] [ ] [ ] ← అన్నీ stopped (DOWNTIME!)
Step 3: [v2] [v2] [v2] ← కొత్త version తో అన్నీ started
Deployment Strategy Comparison
| Strategy | Downtime | Risk | Rollback Speed | Resource Cost | Best For |
|---|---|---|---|---|---|
| Rolling | Zero | Medium | Slow (ఒకదాని తర్వాత ఒకటి roll back) | Low (same infra) | చాలా applications (default) |
| Blue-Green | Zero | Low | Instant (తిరిగి switch) | High (double infra) | Instant rollback కావాల్సిన critical apps |
| Canary | Zero | Very Low | Fast (canary kill చేయండి) | Low-Medium | High-traffic apps, కొత్త risky features |
| Recreate | Yes | High | Slow (full redeploy) | Low | Dev/staging, downtime tolerate చేసే apps |
Observability Stack — Metrics, Logs, Traces, Alerts
Observability మూడు ప్రశ్నలకు సమాధానం ఇస్తుంది: ఏమి జరిగింది? (logs), ఇప్పుడు ఏమి జరుగుతోంది? (metrics), మరియు ఎందుకు జరిగింది? (traces). ఒక doctor లాగా ఆలోచించండి: symptoms (metrics), medical history (logs), మరియు MRI scan (traces).
Metrics: Micrometer → Prometheus → Grafana
ఎలా పని చేస్తుంది:
┌──────────────┐ scrape /metrics ┌────────────┐ dashboards ┌─────────┐
│ Spring Boot │ ────────every 15s───→ │ Prometheus │ ──────────────→ │ Grafana │
│ (Micrometer)│ │ (time-series│ │ (charts)│
└──────────────┘ │ database) │ └─────────┘
// మీ code లో custom metrics
@Service
@RequiredArgsConstructor
public class OrderMetrics {
private final MeterRegistry meterRegistry;
public void recordOrderPlaced(String orderType) {
// Counter — పైకి మాత్రమే వెళ్తుంది (total orders placed)
meterRegistry.counter("orders.placed", "type", orderType).increment();
}
public void recordOrderProcessingTime(long millis) {
// Timer — duration distribution ను track చేస్తుంది
meterRegistry.timer("orders.processing.time")
.record(millis, TimeUnit.MILLISECONDS);
}
public void recordActiveOrders(int count) {
// Gauge — పైకి మరియు కిందకు వెళ్తుంది (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 చేస్తుంది.
ఎలా పని చేస్తుంది:
┌────────────┐ 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 — ప్రతి log line కి ఒక JSON object:
{
"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 ను multiple services గుండా ప్రయాణిస్తున్నప్పుడు follow చేస్తుంది. ప్రతి service ఒక "span" (దాని యాత్రలో భాగం) add చేస్తుంది. సమయం ఎక్కడ ఖర్చు అయిందో మీరు ఖచ్చితంగా చూడవచ్చు.
ఎలా పని చేస్తుంది:
┌──────────┐ ┌──────────┐ ┌──────────┐
│ API │────→│ Order │────→│ Payment │
│ Gateway │ │ Service │ │ Service │
│ span: 2ms│ │ span:45ms│ │ span:30ms│
└──────────┘ └──────────┘ └──────────┘
│ │ │
└────────────────┴────────────────┘
│
▼
┌────────────┐
│ Zipkin │
│ (trace UI) │
└────────────┘
Trace ID: abc-123 (అన్ని services లో same)
Total time: 77ms
Bottleneck: Payment Service (30ms) — investigate!
// application.yml — Tracing configuration
management:
tracing:
sampling:
probability: 1.0 # Dev లో 100% sample, prod లో 10% (0.1)
spring:
application:
name: order-service
// pom.xml dependencies:
// micrometer-tracing-bridge-brave
// zipkin-reporter-brave
Alerting: Prometheus AlertManager → Slack/Email
ఎలా పని చేస్తుంది:
┌────────────┐ 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: "2 నిమిషాలు 10 errors/sec కంటే ఎక్కువ"
- 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 — Alerts ను Slack కి route చేయండి
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
Complete 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 vs Rate Limiter ఎప్పుడు ఉపయోగించాలి?
చాలా ఎక్కువ incoming requests నుండి మీ స్వంత service ను రక్షించడానికి Rate Limiter ఉపయోగించండి (club వద్ద bouncer లాగా). విరిగిపోయిన downstream service ను call చేయడానికి సమయం వృథా చేయకుండా మీ service ను రక్షించడానికి Circuit Breaker ఉపయోగించండి (చనిపోయిన appliance ను unplug చేయడం లాగా, మీ ఇంటి breaker trip కాకుండా). Production లో, రెండూ కలిపి ఉపయోగించండి: incoming ను rate limit చేయండి, outgoing ను circuit break చేయండి.
2. Kafka మరియు Spring Events మధ్య తేడా ఏమిటి?
Spring Events ఒకే application లోపల పని చేస్తాయి — ఒకే room లో ఉన్న వ్యక్తికి note pass చేయడం లాంటిది. Kafka వేర్వేరు servers లో run అవుతున్న multiple applications అంతటా పని చేస్తుంది — post office ద్వారా letter పంపడం లాంటిది. In-process decoupling కోసం Spring Events ఉపయోగించండి (ఉదా., order completion తర్వాత email పంపడం). వేర్వేరు machines లో వేర్వేరు services asynchronously communicate చేయాల్సినప్పుడు, మరియు message durability (messages crashes survive అవుతాయి) అవసరమైనప్పుడు Kafka ఉపయోగించండి.
3. Outbox Pattern data inconsistency ను ఎలా నిరోధిస్తుంది?
Outbox pattern లేకుండా, మీరు database కి save చేసి, తర్వాత Kafka కి publish చేస్తారు — రెండు separate operations. వాటి మధ్య మీ app crash అయితే, database లో data ఉంటుంది కానీ event publish కాలేదు. Outbox pattern business data మరియు event రెండింటినీ same database transaction లో write చేస్తుంది. రెండూ save అవుతాయి లేదా రెండూ save కావు. Separate background process తర్వాత outbox table ను read చేసి Kafka కి publish చేస్తుంది. Publisher crash అయినా, event safely database లో ఉంటుంది, next run లో publish కోసం wait చేస్తూ.
4. Production application కోసం ఏ deployment strategy best?
Rolling update best default choice — దీనికి zero downtime ఉంటుంది, same infrastructure ఉపయోగిస్తుంది, మరియు Kubernetes లో built-in. Instant rollback కావాల్సిన mission-critical apps కోసం, Blue-Green ఉపయోగించండి (కానీ double infrastructure cost ఉంటుంది). కొత్త features safely test చేయాలనుకునే high-traffic apps కోసం, Canary ఉపయోగించండి (మొదట 5% traffic route చేయండి, తర్వాత క్రమంగా పెంచండి). Recreate development environments లేదా కొన్ని నిమిషాల downtime tolerate చేయగల apps కోసం మాత్రమే.
5. Metrics, Logs, మరియు Traces కలిసి ఎలా పని చేస్తాయి?
Production issue ను investigate చేయడం ఒక detective లా ఉంటుంది. Metrics ఏదో తప్పు అని చెప్పే alarm ("2:30 PM కి error rate spike అయింది"). Logs మీకు details ఇస్తాయి ("user-789 కోసం OrderService.java line 42 లో NullPointerException"). Traces services అంతటా ఒక request యొక్క full journey చూపిస్తాయి ("ఈ request order-service లో 200ms ఖర్చు చేసింది, తర్వాత payment-service లో 3000ms stuck — అదే bottleneck!"). కలిసి, అవి సమాధానం ఇస్తాయి: ఏమి జరిగింది, ఎక్కడ జరిగింది, మరియు ఎందుకు జరిగింది.