Kafka advanced with Spring

01.07.2026 | Muallif: Jaxongir a.k.a | Kategoriya: Spring Boot | 27 daqiqa o'qish

Kafka advanced with Spring - bu oddiy event yuborish/qabul qilish emas. Senior darajada siz consumer group, partition, offset, ordering, idempotency, transaction, retry topic, DLT, schema evolution va outbox kabi production muammolarini boshqarishingiz kerak.

Bu blok quyidagilarni o‘z ichiga oladi: consumer group management, exactly-once semantics, Kafka Streams, Schema Registry & Avro/Protobuf, transactional outbox, retry topics, DLT.


1. Kafka’ni oddiy queue deb tushunish xato

RabbitMQ’da ko‘pincha model bunday:

Producer → Queue → Consumer

Kafka’da esa model boshqacha:

Producer → Topic → Partition → Consumer Group → Consumer

Kafka eventlarni topic ichida saqlaydi. Consumer eventni o‘qigandan keyin event topic’dan o‘chib ketmaydi. Kafka rasmiy hujjatlarida eventlar topiclarda tartiblangan va durable saqlanishi, Kafka esa exactly-once processing kabi kafolatlarni qo‘llashi aytiladi. (Apache Kafka)


2. Topic, partition, offset

Topic

Topic - eventlar saqlanadigan nomlangan kanal.

orders.created
payments.completed
notifications.requested

Misol:

order-service → orders.created → inventory-service
                              → payment-service
                              → analytics-service

Partition

Topic bir nechta partition’dan iborat bo‘lishi mumkin.

orders.created
  ├── partition-0
  ├── partition-1
  └── partition-2

Partition - Kafka scaling va ordering mexanizmining yuragi.


Offset

Offset - partition ichidagi event pozitsiyasi.

partition-0:
offset 0 → OrderCreated#1
offset 1 → OrderCreated#2
offset 2 → OrderCreated#3

Consumer qayergacha o‘qiganini offset orqali biladi.


3. Ordering: tartib kafolati

Kafka’da tartib topic bo‘yicha emas, partition bo‘yicha kafolatlanadi.

partition-0 ichida tartib bor
partition-1 ichida tartib bor
lekin partition-0 va partition-1 orasida global tartib yo‘q

Shuning uchun key tanlash juda muhim.

Yomon:

kafkaTemplate.send("orders.created", event);

Bu eventlarni turli partitionlarga randomroq taqsimlaydi.

Yaxshi:

kafkaTemplate.send("orders.created", orderId.toString(), event);

Bu bir xil orderId bo‘yicha eventlar bir partition’ga tushadi.

Order 15 CREATED
Order 15 PAID
Order 15 SHIPPED

hammasi bir partition’da bo‘lsa, tartib saqlanadi.


4. Consumer group management

Consumer group nima?

Bir xil group ichidagi consumerlar topic partitionlarini bo‘lib o‘qiydi.

Topic: orders.created
Partitions: 3

Consumer Group: payment-service
  consumer-1 → partition-0
  consumer-2 → partition-1
  consumer-3 → partition-2

Agar consumer 4 ta bo‘lsa:

consumer-4 → idle

Chunki partition 3 ta. Bir group ichida bitta partition’ni bir vaqtda faqat bitta consumer o‘qiydi.


Boshqa group - mustaqil o‘qish

orders.created topic

payment-service group   → eventlarni o‘qiydi
inventory-service group → eventlarni o‘qiydi
analytics-service group → eventlarni o‘qiydi

Har bir group o‘z offsetiga ega.

payment-service offset: 1200
analytics-service offset: 900
inventory-service offset: 1195

Spring Kafka listener

@Component
@RequiredArgsConstructor
public class OrderCreatedListener {

    private final PaymentService paymentService;

    @KafkaListener(
            topics = "orders.created",
            groupId = "payment-service"
    )
    public void listen(OrderCreatedEvent event) {
        paymentService.startPayment(event.orderId(), event.amount());
    }
}

Bu yerda groupId = "payment-service" juda muhim. Agar groupId noto‘g‘ri tanlansa, eventlar noto‘g‘ri parallelizm yoki duplicate processing bilan ketishi mumkin.


5. Partition soni va consumer soni

Qoida

Parallelism limit = partition soni

Agar topic’da 6 partition bo‘lsa:

max active consumers in same group = 6

Agar 10 consumer bo‘lsa, 4 tasi idle bo‘ladi.


Partition sonini tanlashda savollar

Savol

Sabab

Throughput qancha kerak?

Ko‘proq partition → ko‘proq parallelism

Ordering qayerda kerak?

Bir key bir partition’da qolishi kerak

Consumer processing sekinmi?

Partition ko‘paytirish yordam berishi mumkin

Key distribution tekismi?

Hot partition oldini olish kerak

Kelajakda traffic oshadimi?

Partitionni oldindan rejalash kerak


6. Hot partition muammosi

Agar hamma event bir xil key bilan yuborilsa:

kafkaTemplate.send("orders.created", "ORDER", event);

hammasi bitta partition’ga tushishi mumkin.

partition-0 → 1 000 000 events
partition-1 → 10 events
partition-2 → 15 events

Bu hot partition.

Yaxshi key:

orderId.toString()

Yoki user bo‘yicha tartib kerak bo‘lsa:

userId.toString()

Senior savol:

Menga ordering qaysi entity bo‘yicha kerak?
orderId?
userId?
tenantId?
paymentId?

7. Offset commit

Consumer eventni o‘qigandan keyin offset commit qiladi.

Event o‘qildi
↓
Business logic ishladi
↓
Offset commit qilindi

Agar offset oldin commit qilinsa, lekin business logic fail bo‘lsa:

Event yo‘qolgan hisoblanadi

Agar business logic bajarildi, lekin offset commit bo‘lmasa:

Event qayta o‘qilishi mumkin

Shuning uchun Kafka consumer logic idempotent bo‘lishi kerak.


8. At-most-once, at-least-once, exactly-once

At-most-once

Event 0 yoki 1 marta process bo‘ladi.
Lekin yo‘qolishi mumkin.

Yomon payment/order use-case uchun.


At-least-once

Event kamida 1 marta process bo‘ladi.
Duplicate bo‘lishi mumkin.

Kafka’da ko‘p production systemlar at-least-once + idempotency bilan ishlaydi.


Exactly-once

Event processing natijasi exactly once ko‘rinadi.

Bu eng murakkab. Kafka transactionlari yordam beradi, lekin external DB yoki external API bilan to‘liq exactly-once qilish har doim ham oddiy emas.

Spring Kafka hujjatlarida transactionlardan foydalanish Exactly Once Semantics bilan bog‘lanishi ko‘rsatiladi. (Home)


9. Idempotent consumer

Kafka consumer duplicate event olishi mumkin. Shuning uchun consumer idempotent bo‘lishi kerak.

Muammo

PaymentCompletedEvent ikki marta keldi
↓
Order ikki marta PAID qilindimi?
Notification ikki marta yuborildimi?

Yechim: processed_events table

CREATE TABLE processed_events (
    event_id VARCHAR(100) PRIMARY KEY,
    processed_at TIMESTAMP NOT NULL
);

Listener:

@Component
@RequiredArgsConstructor
public class PaymentCompletedListener {

    private final OrderService orderService;

    @KafkaListener(topics = "payments.completed", groupId = "order-service")
    public void listen(PaymentCompletedEvent event) {
        orderService.markPaid(event);
    }
}

Service:

@Service
@RequiredArgsConstructor
public class OrderService {

    private final ProcessedEventRepository processedEventRepository;
    private final OrderRepository orderRepository;

    @Transactional
    public void markPaid(PaymentCompletedEvent event) {
        if (processedEventRepository.existsById(event.eventId())) {
            return;
        }

        Order order = orderRepository.findById(event.orderId())
                .orElseThrow();

        order.markPaid();

        processedEventRepository.save(
                new ProcessedEvent(event.eventId(), Instant.now())
        );
    }
}

DB unique constraint duplicate’dan himoya qiladi.


10. Producer idempotency

Producer ham duplicate yuborishi mumkin.

Spring Boot config:

spring:
  kafka:
    producer:
      properties:
        enable.idempotence: true
        acks: all

Bu producer retry qilganda broker duplicate yozib qo‘yishini kamaytiradi.

Lekin bu biznes-level idempotency o‘rnini bosmaydi.

Producer idempotency ≠ business idempotency

11. Kafka transaction

Kafka transaction bir nechta Kafka operation’ni atomik qilish uchun ishlatiladi.

Masalan:

Consume topic A
↓
Process
↓
Produce topic B
↓
Offset commit

Bular bitta transaction’da bo‘lishi mumkin.

Agar produce fail bo‘lsa, offset ham commit bo‘lmaydi.

Spring Kafka’da transaction uchun KafkaTransactionManager ishlatiladi.

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

Producer:

@Service
@RequiredArgsConstructor
public class OrderEventPublisher {

    private final KafkaTemplate<String, OrderCreatedEvent> kafkaTemplate;

    @Transactional
    public void publish(OrderCreatedEvent event) {
        kafkaTemplate.send("orders.created", event.orderId().toString(), event);
    }
}

Muhim: bu Kafka transaction. Agar siz DB ham ishlatayotgan bo‘lsangiz, DB transaction bilan Kafka transaction orasidagi atomicity alohida masala.


12. DB + Kafka muammosi

Eng ko‘p uchraydigan production bug:

@Transactional
public void createOrder(CreateOrderRequest request) {
    Order order = orderRepository.save(new Order(request));

    kafkaTemplate.send("orders.created", new OrderCreatedEvent(order.getId()));
}

Muammo:

DB save success
Kafka send fail
↓
Order bor, lekin event yo‘q

Yoki:

Kafka send success
DB commit fail
↓
Event bor, lekin order yo‘q

Bu distributed transaction muammosi.


13. Transactional Outbox pattern

G‘oya

DB va eventni bitta DB transaction ichida saqlaymiz.

orders table
outbox_events table

Service:

@Transactional
public void createOrder(CreateOrderRequest request) {
    Order order = orderRepository.save(new Order(request));

    OutboxEvent event = OutboxEvent.of(
            "OrderCreated",
            order.getId().toString(),
            new OrderCreatedPayload(order.getId(), order.getUserId())
    );

    outboxRepository.save(event);
}

Bu bitta DB transaction:

Order save
Outbox event save
Commit

Keyin background publisher:

outbox_events’dan o‘qiydi
Kafka’ga yuboradi
sent deb belgilaydi

Architecture:

Order Service
  ↓ DB transaction
orders table + outbox_events table
  ↓
Outbox Publisher
  ↓
Kafka topic: orders.created

Bu pattern Kafka bilan Spring microservice’larda juda amaliy.


Outbox table

CREATE TABLE outbox_events (
    id UUID PRIMARY KEY,
    aggregate_type VARCHAR(100) NOT NULL,
    aggregate_id VARCHAR(100) NOT NULL,
    event_type VARCHAR(100) NOT NULL,
    payload JSONB NOT NULL,
    status VARCHAR(30) NOT NULL,
    created_at TIMESTAMP NOT NULL,
    published_at TIMESTAMP
);

Publisher:

@Component
@RequiredArgsConstructor
public class OutboxPublisher {

    private final OutboxRepository outboxRepository;
    private final KafkaTemplate<String, String> kafkaTemplate;

    @Scheduled(fixedDelay = 1000)
    @Transactional
    public void publish() {
        List<OutboxEvent> events = outboxRepository.findTop100ByStatusOrderByCreatedAt("NEW");

        for (OutboxEvent event : events) {
            kafkaTemplate.send(
                    topicFor(event.getEventType()),
                    event.getAggregateId(),
                    event.getPayload()
            );

            event.markPublished();
        }
    }

    private String topicFor(String eventType) {
        return switch (eventType) {
            case "OrderCreated" -> "orders.created";
            default -> throw new IllegalArgumentException("Unknown event type");
        };
    }
}

Senior eslatma: bu kod soddalashtirilgan. Real production’da Kafka send result, retry, lock, batching va duplicate protection alohida ko‘riladi.


14. Debezium bilan Outbox

Outbox publisher’ni o‘zingiz yozmasdan, Debezium CDC orqali ham qilish mumkin.

App writes outbox table
↓
Debezium reads DB transaction log
↓
Kafka topic’ga event chiqaradi

Afzallik:

App Kafka publish bilan shug‘ullanmaydi
DB commit bo‘lgan event transaction log’dan olinadi

Kamchilik:

Infra murakkabroq
Debezium, connector, schema, monitoring kerak

15. Retry topics

Consumer eventni process qila olmadi.

Sabablar:

- External service timeout
- DB deadlock
- Temporary network issue
- Downstream service 503

Bunday holatda darhol qayta-qayta retry qilish yomon:

Main consumer thread bloklanadi
Kafka lag oshadi
Boshqa eventlar kutadi

Spring Kafka’da non-blocking retry uchun retry topic modeli ishlatiladi. Spring Kafka hujjatlarida @RetryableTopic va DLT processing uchun @DltHandler ishlatilishi ko‘rsatilgan. (Home)


@RetryableTopic misol

@Component
@RequiredArgsConstructor
public class OrderEventListener {

    private final OrderProcessor orderProcessor;

    @RetryableTopic(
            attempts = "4",
            backoff = @Backoff(delay = 1000, multiplier = 2.0),
            dltTopicSuffix = ".dlt"
    )
    @KafkaListener(topics = "orders.created", groupId = "inventory-service")
    public void listen(OrderCreatedEvent event) {
        orderProcessor.process(event);
    }

    @DltHandler
    public void handleDlt(OrderCreatedEvent event, Exception ex) {
        log.error("Order event moved to DLT. orderId={}, error={}",
                event.orderId(),
                ex.getMessage()
        );
    }
}

Flow:

orders.created
  ↓ fail
orders.created-retry-1000
  ↓ fail
orders.created-retry-2000
  ↓ fail
orders.created-retry-4000
  ↓ fail
orders.created.dlt

16. DLT - Dead Letter Topic

DLT - qayta ishlanmagan xato eventlar tushadigan topic.

Confluent DLT’ni downstream consumer process qila olmagan message’lar yuboriladigan maxsus Kafka topic sifatida tushuntiradi; bu pipeline ishlashda davom etishiga va operatorlarga xato message’larni tekshirishga yordam beradi. (Confluent)

orders.created.dlt
payments.completed.dlt
notifications.requested.dlt

DLT’ga tushgan eventni:

- inspect qilish
- alert qilish
- manual fix qilish
- reprocess qilish
- discard qilish

mumkin.


DLT’da nima saqlash kerak?

DLT message’da kerakli metadata bo‘lsin:

original topic
partition
offset
consumer group
exception class
exception message
timestamp
traceId
eventId
payload

Bu debugging uchun juda muhim.


17. Retry qilish mumkin bo‘lgan va bo‘lmagan xatolar

Xato

Retry?

Sabab

Network timeout

Ha

Vaqtinchalik bo‘lishi mumkin

HTTP 503

Ha

Downstream vaqtincha down

DB deadlock

Ha

Qayta urinishda o‘tishi mumkin

Validation error

Yo‘q

Event noto‘g‘ri

Unknown enum

Yo‘q

Schema/contract muammo

Missing required field

Yo‘q

Data muammo

Permission denied

Yo‘q

Retry foydasiz

Duplicate event

Yo‘q

Idempotent skip qilish kerak

Senior qoida:

Har xatoga retry qilmang.
Retry faqat vaqtinchalik xatolar uchun.

18. Schema Registry

Muammo

Producer event structure’ni o‘zgartirdi.

Oldingi event:

{
  "orderId": 10,
  "amount": 100000
}

Yangi event:

{
  "id": 10,
  "totalAmount": 100000,
  "currency": "UZS"
}

Consumer eski format kutayotgan bo‘lsa, yiqiladi.


Schema Registry nima qiladi?

Schema Registry event schema’larini markaziy joyda saqlaydi va compatibility qoidalarini tekshiradi. Confluent Schema Registry Avro, JSON Schema va Protobuf serializer/deserializer’larini qo‘llashi, schema’lar tarixini saqlashi va compatibility settings orqali schema evolution’ga ruxsat berishini hujjatlarda ko‘rsatadi. (Confluent Documentation)

Producer → Schema Registry’dan schema ID oladi
Producer → Kafka’ga schema ID + payload yuboradi
Consumer → schema ID orqali decode qiladi

19. Avro / Protobuf / JSON Schema

Format

Kuchli tomoni

Kamchiligi

JSON

Oddiy, o‘qish oson

Contract zaif, hajm katta

Avro

Kafka ecosystem’da juda keng, schema evolution kuchli

Human-readable emas

Protobuf

Compact, tez, cross-language yaxshi

Field number discipline kerak

JSON Schema

JSON bilan contract

Hajm kattaroq


Avro schema misol

{
  "type": "record",
  "name": "OrderCreatedEvent",
  "namespace": "uz.company.orders",
  "fields": [
    { "name": "eventId", "type": "string" },
    { "name": "orderId", "type": "long" },
    { "name": "userId", "type": "long" },
    { "name": "amount", "type": "double" },
    { "name": "currency", "type": "string", "default": "UZS" }
  ]
}

Yangi field qo‘shsangiz, default qiymat berish muhim:

{ "name": "currency", "type": "string", "default": "UZS" }

20. Schema evolution

Confluent hujjatlarida Schema Registry default compatibility type BACKWARD ekanligi ko‘rsatilgan; u Avro, Protobuf va JSON Schema uchun compatibility qoidalarini qo‘llaydi. (Confluent Documentation)

Backward compatible o‘zgarishlar

Odatda xavfsizroq:

- Optional field qo‘shish
- Default value bilan field qo‘shish
- Eski consumer o‘qiy oladigan formatni saqlash

Breaking change

Xavfli:

- Field nomini o‘zgartirish
- Field type’ni o‘zgartirish
- Required field qo‘shish
- Eski fieldni birdan o‘chirib tashlash

Yomon:

"orderId": "long"

dan:

"orderId": "string"

ga o‘tish. Consumerlar buzilishi mumkin.


21. Kafka Streams with Spring

Kafka Streams - Kafka ustida stream processing qilish uchun library.

Masalan:

orders.created + payments.completed → order-status-updated

Spring Kafka hujjatlarida Spring for Apache Kafka Kafka Streams uchun first-class support berishi aytiladi; ishlatish uchun kafka-streams jar classpath’da bo‘lishi kerak. (Home)

Confluent Kafka Streams’da stream "unbounded, continuously updating data set" sifatida tushuntiriladi. (Confluent Documentation)


Kafka Streams qachon kerak?

- Eventlarni join qilish
- Aggregation qilish
- Windowing
- Real-time metrics
- Fraud detection
- Event enrichment
- Stream’dan stream hosil qilish

KStream misol

@Configuration
public class OrderStreamConfig {

    @Bean
    public KStream<String, OrderCreatedEvent> orderStream(StreamsBuilder builder) {
        KStream<String, OrderCreatedEvent> orders =
                builder.stream("orders.created");

        orders
                .filter((key, order) -> order.amount().compareTo(BigDecimal.ZERO) > 0)
                .mapValues(order -> new ValidOrderEvent(
                        order.eventId(),
                        order.orderId(),
                        order.amount()
                ))
                .to("orders.validated");

        return orders;
    }
}

Bu listener emas. Bu stream topology.


22. KStream vs KTable

Tushuncha

Ma’nosi

KStream

Event oqimi

KTable

Oxirgi holat / changelog table

GlobalKTable

Hamma instance’da to‘liq table copy

Misol:

OrderCreatedEvent → KStream
UserProfileUpdated → KTable

Agar order eventini user ma’lumoti bilan boyitmoqchi bo‘lsangiz:

orders KStream join users KTable

23. Event design

Yomon event:

{
  "data": "something happened"
}

Yaxshi event:

{
  "eventId": "evt-123",
  "eventType": "OrderCreated",
  "eventVersion": 1,
  "occurredAt": "2026-05-21T10:00:00Z",
  "tenantId": "company-12",
  "traceId": "abc-123",
  "payload": {
    "orderId": 100,
    "userId": 55,
    "amount": 120000,
    "currency": "UZS"
  }
}

Senior event’da quyidagilar bo‘lishi kerak:

eventId
eventType
eventVersion
occurredAt
aggregateId
tenantId
traceId/correlationId
payload

24. Topic naming

Yaxshi nomlash:

orders.created
orders.cancelled
payments.completed
payments.failed
notifications.requested

Yomon:

topic1
events
order
test-topic

Topic nomi business eventni bildirishi kerak.


25. Kafka monitoring

Kafka production’da observability juda muhim.

Kuzatiladigan metrikalar:

consumer lag
records consumed/sec
records produced/sec
request latency
error rate
retry topic size
DLT message count
rebalance count
under-replicated partitions
broker disk usage

Spring Boot + Micrometer bilan Kafka listener metrics va custom counters qo‘shish mumkin.

Masalan:

@Component
@RequiredArgsConstructor
public class DltMetrics {

    private final MeterRegistry meterRegistry;

    public void incrementDlt(String topic) {
        meterRegistry.counter("kafka.dlt.messages", "topic", topic).increment();
    }
}

26. Rebalancing

Consumer group’da consumer qo‘shilsa yoki o‘chsa, Kafka partition assignment’ni qayta taqsimlaydi.

consumer-1 down
↓
partitionlar consumer-2 va consumer-3 ga qayta beriladi

Bu rebalance.

Muammo:

Rebalance paytida processing pauza bo‘lishi mumkin

Sabablar:

- Consumer juda sekin
- max.poll.interval.ms oshib ketadi
- Processing listener ichida juda uzoq
- Consumer crash bo‘ladi

Yechimlar:

- Processingni tezlatish
- Batch o‘lchamini to‘g‘ri tanlash
- max.poll.interval.ms sozlash
- Heavy workni alohida workerga berish
- Idempotency qo‘shish

27. Batch listener

Agar throughput kerak bo‘lsa, batch listener ishlatish mumkin.

@KafkaListener(
        topics = "analytics.events",
        groupId = "analytics-service",
        containerFactory = "batchFactory"
)
public void listen(List<AnalyticsEvent> events) {
    analyticsService.processBatch(events);
}

Factory:

@Bean
ConcurrentKafkaListenerContainerFactory<String, AnalyticsEvent> batchFactory(
        ConsumerFactory<String, AnalyticsEvent> consumerFactory
) {
    ConcurrentKafkaListenerContainerFactory<String, AnalyticsEvent> factory =
            new ConcurrentKafkaListenerContainerFactory<>();

    factory.setConsumerFactory(consumerFactory);
    factory.setBatchListener(true);

    return factory;
}

Batch yaxshi:

- Analytics
- Bulk insert
- Logging pipeline
- Metrics ingestion

Batch ehtiyot:

- Payment
- Order state transition
- Strict per-event error handling kerak bo‘lsa

28. Error handling: blocking retry vs non-blocking retry

Blocking retry

Consumer shu eventda kutib qoladi:

event fail
↓
sleep
↓
retry
↓
sleep
↓
retry

Kam traffic’da mumkin, lekin katta topic’da lag oshiradi.

Non-blocking retry

Event retry topic’ga ketadi:

main topic → retry topic → retry topic → DLT

Bu main consumer’ni bloklamaydi.

Spring Kafka’da retry topic va DLT listener container’larini alohida sozlash imkoniyatlari mavjud. (Home)


29. Real production flow

Order yaratish:

Client
  ↓
Order Service
  ↓ DB transaction
orders table + outbox_events table
  ↓
Outbox Publisher / Debezium
  ↓
Kafka: orders.created
  ├── payment-service group
  ├── inventory-service group
  ├── notification-service group
  └── analytics-service group

Payment service:

orders.created
  ↓
payment-service consumer
  ↓
idempotency check
  ↓
external payment provider
  ↓
payments.completed / payments.failed

Order service:

payments.completed
  ↓
order-service consumer
  ↓
processed_events check
  ↓
order status = PAID

30. Common mistakes

1. Kafka’ni HTTP o‘rniga ishlatish

Yomon:

User request yubordi
↓
Kafka’ga event
↓
Consumer javob qaytarishini kutish

Kafka async event-driven flow uchun. Request/response uchun har doim ham mos emas.


2. Event ichida juda katta payload

Yomon:

Full order
Full user profile
Full product list
Full binary file

Yaxshi:

eventId
orderId
minimal business data

Katta data kerak bo‘lsa, object storage link yoki DB reference ishlatiladi.


3. Retry hamma error uchun

Validation error retry qilinsa, DLT’ga tushguncha resurs yeb qo‘yadi.


4. Idempotency yo‘q

Kafka’da duplicate bo‘lishi mumkin deb loyihalash kerak.


5. Schema version yo‘q

Event contract o‘zgarganda consumerlar buziladi.


6. Topic partition key noto‘g‘ri

Ordering buziladi yoki hot partition paydo bo‘ladi.


31. Senior interview savol-javob

Suhbatda so‘rasa:

Kafka’ni Spring Boot’da production’da qanday ishlatasiz?

Javob:

Kafka’da avval topic, partition, key va consumer group modelini to‘g‘ri loyihalayman.
Ordering kerak bo‘lsa, event key’ni aggregateId, masalan orderId bo‘yicha tanlayman.
Consumer group parallelism partition soni bilan cheklanishini hisobga olaman.

Consumerlar at-least-once ishlashi mumkin, shuning uchun idempotent consumer qilaman:
eventId yoki business key bo‘yicha processed_events table yoki unique constraint ishlataman.

DB bilan Kafka publish qilishda dual-write muammosi bor, shuning uchun transactional outbox yoki Debezium CDC ishlataman.
Temporary xatolar uchun retry topic, doimiy xatolar uchun DLT ishlataman.
DLT monitoring va reprocessing jarayoni bo‘lishi kerak.

Event contract uchun Schema Registry, Avro yoki Protobuf ishlataman.
Breaking change qilmaslik uchun schema evolution va compatibility rule’larni nazorat qilaman.

Exactly-once Kafka ichida transaction bilan mumkin, lekin external DB/API bilan ishlaganda idempotency va outbox baribir muhim.

32. Qisqa xulosa

Kafka advanced’da eng muhim qoida:

Kafka duplicate event yuborishi yoki consumer qayta process qilishi mumkin deb loyihala.

Asosiy tushunchalar:

Mavzu

Esda qoladigan qoida

Partition

Ordering partition ichida

Key

Bir aggregate eventlarini bir partition’ga tushiradi

Consumer group

Har group mustaqil o‘qiydi

Offset

Qayergacha o‘qilganini bildiradi

Idempotent consumer

Duplicate’dan himoya

Exactly-once

Kafka ichida mumkin, external system bilan murakkab

Outbox

DB + Kafka dual-write muammosini kamaytiradi

Retry topic

Main consumer’ni bloklamasdan retry

DLT

Tuzalmaydigan eventlar uchun xavfsiz joy

Schema Registry

Event contract’ni boshqaradi

Kafka Streams

Real-time stream processing