Bu qatlamda nima bor
Core darslik mental model va asosiy production muammolarini berdi. Bu Deep qatlam senior darajadagi chuqur mexanikani ochadi: yetkazish semantikasi, tartib, transactional messaging (CDC), retry/DLQ strategiyasi, Kafka ichki ishlashi, sxema kontraktlari, kuzatuvchanlik va ishlaydigan test. Oxirida — senior tayyorlik cheklisti.
At-least-once, exactly-once va idempotentlik
inbox · EOSSenior haqiqat: "exactly-once delivery" — deyarli afsona. Amalda biz "exactly-once effect"ga erishamiz: at-least-once yetkazish + idempotent qabul qiluvchi. Bu farqni tushunish — senior va middle orasidagi chegara.
Idempotent producer + idempotent consumer
Ikki tomon ham himoyalanadi. Producer tomon: har xabarga barqaror eventId beriladi (retry’da o‘zgarmaydi). Consumer tomon — Inbox pattern: qabul qiluvchi ishlangan eventIdlarni bazada saqlaydi va yon ta’sirni o‘sha tranzaksiyada bajaradi. Shunda xabar 5 marta kelsa ham, ish bir marta bajariladi:
// INBOX (idempotent receiver): yon ta'sir + "ishlandi" belgisi BITTA tranzaksiyada
async handle(data) {
await this.prisma.$transaction(async (tx) => {
const seen = await tx.inbox.findUnique({ where: { id: data.eventId } });
if (seen) return; // takror -> hech narsa qilma
await tx.inbox.create({ data: { id: data.eventId } }); // belgilaymiz
await tx.notification.create({ data: { ... } }); // asl ish
});
// tranzaksiya commit bo'lgach ack qilamiz
}
Core darsda "processed table"ni ko‘rdik. Inbox’ning farqi — tekshiruv + yon ta’sir + belgilash bitta atomik tranzaksiyada. Aks holda: ishni bajardingu, belgilashdan oldin qulasang — qayta ishlanadi (idempotentlik buziladi).
Kafka EOS (exactly-once semantics) — chegaralari bilan
Kafka haqiqiy "exactly-once"ni faqat o‘z ichida (Kafka’dan o‘qib → ishlab → Kafka’ga yozish) transaksiyalar va idempotent producer orqali beradi. Tashqi tizimga (Postgres, email) bu kafolat cho‘zilmaydi:
// Kafka EOS (exactly-once) — FAQAT Kafka ICHIDA (read -> process -> write atomik)
const producer = kafka.producer({ idempotent: true }); // dedup id
const consumer = kafka.consumer({ groupId: 'orders' });
await producer.transaction(async (txn) => {
await txn.send({ topic: 'order.enriched', messages: [...] }); // yozuv
await txn.sendOffsets({ consumerGroupId: 'orders', topics }); // offset birga commit
});
// MUHIM: bu tashqi DB/email'gacha cho'zilmaydi -> tashqi yon ta'sir uchun baribir
// at-least-once + idempotent receiver (inbox) kerak.
Tizimingni at-least-once deb loyihalab, har bir consumerni idempotent qil (inbox yoki tabiiy idempotentlik). "Exactly-once"ni qidirib vaqt sarflama — yuqoridagi kombinatsiya production’da yetarli va to‘g‘ri.
Tartib — partition, single-writer, retry
per-key orderingSenior haqiqat: global tartib (hamma eventlar bitta qatorda) — qimmat va deyarli kerak emas. Sen‘ga kerak bo‘lgani — per-key (aggregate bo‘yicha) tartib: bitta hisob/buyurtma bo‘yicha eventlar to‘g‘ri ketma-ketlikda kelsin.
Tartibni qanday saqlaymiz
Kafka: bir xil key doim bir partisiyaga tushadi, partisiya ichida tartib kafolatlangan. RabbitMQ: consistent-hash exchange yoki single active consumer bilan per-key tartib. Asosi — single-writer prinsipi: bitta kalitni bir vaqtda faqat bitta consumer ishlasin.
// PER-KEY ordering: bir aggregate (userId) bo'yicha tartib saqlanishi shart
this.kafka.emit('order.events', {
key: order.userId, // SHU key -> doim BITTA partisiya -> tartiblangan
value: { eventId, type: 'order.created', orderId: order.id },
});
// userId=42 ning barcha eventlari bitta partisiyada, ketma-ket ishlanadi.
// Boshqa userlar boshqa partisiyalarda parallel ketadi (masshtab + tartib birga).
msg1 xato berib retry’ga tushdi, ayni paytda msg2 muvaffaqiyatli ishlandi — endi msg2 msg1dan oldin bajarildi, tartib buzildi. Yechimlar: (1) kalit bo‘yicha xatoda to‘xtab, retry tugamaguncha keyingisini ishlamaslik; (2) handlerni versiyaga chidamli qilish — har aggregate’da version saqlab, eski (kech kelgan) eventni rad etish.
Tartib chinakam kerakmi, o‘ylab ko‘r. Ko‘p hollarda handlerni kommutativ/idempotent qilish (tartibga bog‘liq bo‘lmaslik) — qattiq tartib kafolatidan ko‘ra arzon va mustahkam.
Outbox → CDC (Debezium) + Inbox
transactional messagingCore darsda polling outbox (Cron) ko‘rdik — bu ishlaydi, lekin productionda kamchiliklari bor: qo‘shimcha DB yuki, kechikish (poll oralig‘i), tartib va tozalash muammolari.
Production usuli — CDC (Change Data Capture) / Debezium
Eng yetuk yechim: ma’lumotlar bazasining tranzaksiya jurnalini (WAL) kuzatish. Debezium Postgres WAL’ni o‘qiydi va outbox jadvaliga tushgan har yozuvni avtomatik Kafka’ga uzatadi — poller yo‘q, kechikish minimal, tartib jurnaldan tabiiy keladi.
// prisma/schema.prisma — outbox va inbox jadvallari
model Outbox {
id String @id @default(uuid())
aggregate String // "order"
type String // "order.created"
payload Json
createdAt DateTime @default(now())
@@index([createdAt])
}
model Inbox { // consumer tomonda idempotentlik
id String @id // eventId
consumedAt DateTime @default(now())
}// PRODUCTION usuli: Debezium CDC — Postgres WAL'ni o'qib, outbox'ni Kafka'ga uzatadi
// connector config (JSON) — poller YO'Q, kechikish past, tartib WAL'dan tabiiy
{
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "postgres",
"table.include.list": "public.outbox",
"transforms": "outbox",
"transforms.outbox.type":
"io.debezium.transforms.outbox.EventRouter", // outbox -> topic router
"transforms.outbox.route.by.field": "type" // type -> topic nomi
}
| Polling outbox | CDC (Debezium) | |
|---|---|---|
| Kechikish | Poll oralig‘icha (sekundlar) | Deyarli darhol |
| DB yuki | Doimiy so‘rov | WAL o‘qish (arzon) |
| Murakkablik | Past — oddiy Cron | Yuqori — Kafka Connect + Debezium |
| Mos | Kichik/o‘rta hajm | Katta hajm, qat’iy kafolat |
Inbox (1-bo‘lim) — consumer tomonda dual-write’ning teskari muammosini hal qiladi. "Listen to yourself" — producer o‘z eventini o‘zi ham tinglab, o‘qish modelini yangilaydi; shunda yozuv yo‘li va o‘qish yo‘li bitta event oqimi orqali izchil bo‘ladi (CQRS bilan bog‘lanadi — 9-mavzu).
Retry, backoff, DLQ va parking lot
delayed retry · replaySenior haqiqat: xato bo‘lgan xabarni darhol cheksiz qayta navbatga solish (requeue) — eng ko‘p uchraydigan xato. Bu "poison message" butun navbatni bloklaydi va tizimni yiqitadi.
To‘g‘ri retry strategiyasi
Kechiktirilgan retry (delayed retry): xato bo‘lsa, darhol emas — kutib qayta urinish, har safar oraliqni oshirib. RabbitMQda buni TTL + DLX bilan "retry navbatlari" orqali, Kafka’da esa alohida "retry topic"lar orqali quriladi.
// RabbitMQ kechiktirilgan retry: asosiy navbat -> retry navbati (TTL) -> qaytadi
// retry_queue arguments:
{
'x-message-ttl': 5000, // 5s kut
'x-dead-letter-exchange': 'main.exchange', // TTL tugagach asosiyga qaytar
'x-dead-letter-routing-key': 'order.created',
}
@EventPattern('order.created')
async handle(@Payload() data, @Ctx() ctx) {
const ch = ctx.getChannelRef(); const msg = ctx.getMessage();
const attempts = (msg.properties.headers['x-death']?.[0]?.count) ?? 0;
try {
await this.work(data);
ch.ack(msg);
} catch (e) {
if (attempts >= 5) {
ch.nack(msg, false, false); // 5 urinish bo'ldi -> DLQ ("parking lot")
} else {
ch.publish('retry.exchange', 'order.created', msg.content); // retry navbatiga
ch.ack(msg);
}
}
}
Parking lot (DLQ) + replay
Belgilangan urinishdan keyin ham bo‘lmasa — xabar DLQga ("parking lot") suriladi, navbat bo‘shaydi. DLQ — chetga tashlangan joy emas, operatsion vosita: u monitoring ostida bo‘lishi va undan replay (qayta o‘ynatish) imkoniyati bo‘lishi shart. Idempotentlik (1-bo‘lim) replay’ni xavfsiz qiladi.
Kafka’da retry’ni asosiy oqimni bloklamasdan qilish muhim: xato xabarni retry-topicga ko‘chir, asosiy partisiya ishlashda davom etsin. Aks holda bitta sekin xabar butun partisiyani to‘xtatadi.
Consumer group, partition, rebalancing, lag
kafka internalsKafka’ni "shunchaki broker" deb bilish — middle daraja. Senior partisiya, consumer group, rebalancing va offsetni tushunadi — chunki bularsiz tartib, masshtab va lag’ni boshqarib bo‘lmaydi.
4 ta yadro tushuncha
- Partition — topic’ning bo‘lagi; parallellik va tartib birligi. Tartib faqat partisiya ichida.
- Consumer group — har partisiyani guruh ichidagi aniq bitta consumer o‘qiydi. Demak parallellik darajasi ≤ partisiyalar soni (ortiqcha consumer bekor turadi).
- Offset — consumer qayergacha o‘qiganini bildiruvchi belgi. Ishlangach commit qilsang → at-least-once; oldin commit qilsang → at-most-once.
- Consumer lag = oxirgi offset − commit qilingan offset. Tizim sog‘lig‘ining №1 metrikasi — o‘sib borsa, consumer yetishmayapti.
// NestJS Kafka consumer — CONSUMER GROUP va parallel ishlash
const app = await NestFactory.createMicroservice<MicroserviceOptions>(AppModule, {
transport: Transport.KAFKA,
options: {
client: { brokers: ['localhost:9092'] },
consumer: {
groupId: 'inventory-group', // bir guruh -> partisiyalar taqsimlanadi
sessionTimeout: 30000,
},
run: {
autoCommit: false, // QO'LDA commit -> at-least-once nazorati
partitionsConsumedConcurrently: 3, // 3 partisiyani parallel
},
},
});
Consumer qo‘shilsa/ketsa, partisiyalar qayta taqsimlanadi. Klassik (eager) rebalancing — "stop-the-world": barcha consumerlar to‘xtaydi. Yechim: cooperative-sticky (bosqichma-bosqich) strategiya va uzun ishlov bermaslik (max.poll.interval.msdan oshmaslik), aks holda consumer "o‘lik" deb hisoblanib, doimiy rebalancing bo‘ladi.
Partisiya sonini keyin oshirish mumkin, lekin kamaytirib bo‘lmaydi — va oshirish key → partition moslamasini buzadi (tartib kafolati buziladi). Shuning uchun partisiya sonini oldindan, kelajakdagi parallellikni hisobga olib tanlang.
Sxema evolyutsiyasi va Schema Registry
compatibility · tolerant readerSenior haqiqat: event — bu sen bilmagan ko‘plab consumer bog‘langan abadiy ommaviy kontrakt. Uni o‘zgartirish — API’ni buzish bilan teng. Buni boshqarish — alohida intizom.
Schema Registry va moslik (compatibility) rejimlari
Katta tizimlarda sxema markazlashgan Schema Registryda (Avro/Protobuf/JSON Schema) saqlanadi va har o‘zgarish moslikka tekshiriladi:
| Rejim | Ma’nosi | Kim avval yangilanadi |
|---|---|---|
| BACKWARD | Yangi sxema eski ma’lumotni o‘qiy oladi | Consumer avval |
| FORWARD | Eski sxema yangi ma’lumotni o‘qiy oladi | Producer avval |
| FULL | Ikkala yo‘nalish ham | Tartib muhim emas |
// Avro sxema + Schema Registry (Confluent) — kontraktni markazda boshqarish
{
"type": "record", "name": "OrderCreated",
"fields": [
{ "name": "eventId", "type": "string" },
{ "name": "orderId", "type": "string" },
{ "name": "total", "type": "double" },
{ "name": "currency","type": ["null","string"], "default": null } // YANGI: optional
]
}
// Registry compatibility = BACKWARD -> consumerlar avval yangilanadi.
Tolerant Reader pattern
Consumer sxemaga qat’iy bog‘lanmasin: faqat kerakli maydonlarni o‘qisin, noma’lum maydonlarni e’tiborsiz qoldirsin, yo‘q maydonga default bersin. Shunda producer additive o‘zgarish kiritsa, consumer sinmaydi:
// TOLERANT READER: noma'lum maydonlarni e'tiborsiz qoldir, yo'qiga default ber
function parseOrderCreated(raw: any): OrderCreated {
return {
eventId: raw.eventId,
orderId: raw.orderId,
total: raw.total,
currency: raw.currency ?? 'UZS', // eski eventda yo'q -> default
// raw'dagi boshqa/yangi maydonlar -> shunchaki e'tiborsiz (xato bermaydi)
};
}
Consumer-driven contract test: har consumer "menga shu maydonlar kerak" deb kontrakt yozadi; CI producer o‘zgarishi shu kontraktni buzmasligini avtomatik tekshiradi — eventni sindirmasdan deploy qilasan. Upcasting (event sourcing’da, 10-mavzu): eski versiyadagi eventlarni o‘qishda yangi formatga "ko‘taradigan" funksiya.
Tracing, metrikalar va backpressure
OpenTelemetry · prefetchEDAda xato asinxron, bir necha hop narida sodir bo‘ladi. Kuzatuvchanliksiz (12-mavzu) — bu qorong‘i quti. Senior tizimni tashqaridan tushunarli qilib quradi.
Distributed tracing — async hop’lar bo‘ylab
Sirq: trace kontekstini (traceparent, W3C standart) xabar header’i orqali uzatish. OpenTelemetry produce’da inject, consume’da extract qiladi — natijada butun oqim (HTTP → Orders → broker → Notifications) bitta trace bo‘lib ko‘rinadi:
// DISTRIBUTED TRACING: trace kontekstini xabar header'i orqali uzatamiz
import { propagation, context, trace } from '@opentelemetry/api';
// PRODUCE — joriy trace'ni header'ga "inject" qilamiz
const headers: Record<string, string> = {};
propagation.inject(context.active(), headers);
this.bus.emit('order.created', { value: payload, headers });
// CONSUME — header'dan "extract" qilib, span ochamiz (oqim uzilmaydi)
const parentCtx = propagation.extract(context.active(), msg.properties.headers);
const span = trace.getTracer('notifications')
.startSpan('handle order.created', {}, parentCtx);
// endi Jaeger'da: HTTP -> Orders -> broker -> Notifications BITTA trace
Kuzatiladigan metrikalar
- Consumer lag — №1 metrika; o‘ssa, consumer yetishmayapti (alert qo‘y).
- Processing latency va error rate — har consumer bo‘yicha.
- DLQ depth — DLQ’ga xabar tushsa darhol xabardor bo‘l.
- Throughput (msg/sek) — quvvat rejalashtirish uchun.
Backpressure — consumerni cho‘ktirmaslik
Producer consumerdan tez bo‘lsa, xabarlar to‘planadi. Consumer "hammasini birdan" olib cho‘kib qolmasligi uchun prefetch (QoS) bilan oqim cheklanadi:
// BACKPRESSURE (RabbitMQ): bir consumer bir vaqtda nechta unacked xabar ushlaydi
options: {
// ...
prefetchCount: 10, // 10 dan ortiq tasdiqlanmagan xabar olmaydi -> consumer cho'kmaydi
noAck: false,
}
// Kafka muqobili: max.poll.records + kerak bo'lsa pause()/resume()
Ishlaydigan lab, test va senior cheklist
testcontainersNazariyani ishlaydigan loyihada mustahkamlash — senior bo‘lishning yagona yo‘li. Quyida lab’ning yadro qismlari.
1 · Infratuzilma
# docker-compose.yml — to'liq lab uchun
services:
rabbitmq:
image: rabbitmq:3-management
ports: ["5672:5672", "15672:15672"]
postgres:
image: postgres:16
environment: { POSTGRES_PASSWORD: dev, POSTGRES_DB: shop }
ports: ["5432:5432"]
2 · Integration test (testcontainers) — idempotentlikni isbotlash
Senior darajaning belgisi: EDAni test bilan kafolatlash. Testcontainers real broker ko‘taradi, sen esa "bir event 2 marta kelsa, ish 1 marta bajariladi"ni avtomatik tekshirasan:
// Integration test — testcontainers bilan REAL RabbitMQ ko'tarib tekshirish
import { RabbitMQContainer } from '@testcontainers/rabbitmq';
it('order.created -> email yuboriladi', async () => {
const rmq = await new RabbitMQContainer().start();
// ... app'ni rmq.getAmqpUrl() bilan ulaymiz
await producer.emit('order.created', { eventId: 'e1', orderId: 'o1' });
await waitFor(() => expect(emailSpy).toHaveBeenCalledWith('o1'));
// idempotentlik: aynan shu eventni qayta yuboramiz
await producer.emit('order.created', { eventId: 'e1', orderId: 'o1' });
await waitFor(() => expect(emailSpy).toHaveBeenCalledTimes(1)); // 1 marta!
await rmq.stop();
});
3 · Senior tayyorlik cheklisti
Quyidagilarning hammasiga "ha" deb javob bera olsang — bu mavzuda senior darajadasan:
| Savol / amaliyot | Bog‘liq bo‘lim |
|---|---|
| at-least-once + idempotent receiver (inbox) nima uchun "exactly-once effect" berishini tushuntira olaman | 1 |
| Per-key tartibni partisiya/consistent-hash bilan ta’minlay olaman; retry tartibni buzishini bilaman | 2 |
| Dual-write muammosini Outbox bilan, productionda CDC (Debezium) bilan hal qilaman | 3 |
| Cheksiz requeue qilmayman; backoff retry + DLQ + replay quraman | 4 |
| Consumer group, partition, offset commit va consumer lag’ni boshqaraman | 5 |
| Sxemani BACKWARD/FORWARD compatibility va tolerant reader bilan rivojlantiraman | 6 |
| Trace kontekstini header orqali uzatib, async oqimni bitta trace’da ko‘ raman | 7 |
| Backpressure (prefetch/QoS) sozlayman; EDAni testcontainers bilan testlayman | 7, 8 |
Bu 3-mavzuning Deep/Senior qatlami edi. Core darslik bilan birga — endi bu mavzu junior’dan senior gacha to‘liq.
Tayyor bo‘lsangiz: “4-mavzu: Caching” — Core’dan boshlaymiz, kerak bo‘lsa Deep qatlamini ham qo‘shamiz.