Dilshod.dev

Qulashdan omon chiqadigan background job'lar

Deploy worker'larni ish o'rtasida qayta ishga tushirdi va xatlar shunchaki yo'q bo'ldi. O'lgan jarayon ma'lumot yo'qotish emas, shunchaki kechikish bo'ladigan queue qanday quriladi.

Muallif: Dilshod Abdullayev11 daqiqa o'qish

Payshanba kuni tushdan keyin deploy qildik. Rolling deploy edi, ikkita worker ketma-ket almashtirildi, health check'lar yashil, error tracker'da bir og'iz shikoyat yo'q. Juma kuni ertalab hamkasbim so'radi: kecha ro'yxatdan o'tgan ikki yuzga yaqin foydalanuvchi nega welcome xatini olmadi?

Ko'rib chiqadigan exception yo'q edi. Job'lar olingan, jarayon yuborish o'rtasida SIGTERM olgan, Kubernetes esa SIGKILLdan oldin odatdagi grace period'ni bergan. Ish o'sha oraliqda yo'qoldi. Queue nuqtai nazaridan bu job'lar tarqatilgan va qaytib kelmagan; bizning nuqtai nazarimizdan esa hech narsa buzilmagan, chunki hech narsa buzilganini aytmagan.

Men duch kelgan background job hodisalarining aksariyati aynan shu shaklda bo'ladi. Stack trace'da ko'rinadigan qulash emas — mavjud bo'lishi kerak bo'lgan narsalarni sanab ko'rgandagina sezadigan qulash.

Tugatolmaydigan ishni o'ziga olgan request handler

Xato asl ko'rinishida mana bunday, va u deyarli har bir kodbazada bir paytlar uchraydi:

@Post('signup')
async signup(@Body() dto: SignupDto) {
  const user = await this.users.create(dto);
 
  await this.mailer.sendWelcome(user.email);   // 400ms, ba'zan 9s, ba'zan throw qiladi
  await this.crm.syncContact(user);            // uchinchi tomon, SLA sizning qo'lingizda emas
 
  return { id: user.id };
}

Foydalanuvchining HTTP so'rovi endi unga javob berishga hech qanday aloqasi bo'lmagan ikkita ish uchun javobgar. Pochta provayderi sekin bo'lsa — signup sekin. CRM o'chib qolsa — signup 500 qaytaradi, holbuki foydalanuvchi allaqachon yaratilgan. U qayta uradi va endi sizda xat muammosining ustiga dublikat akkaunt muammosi qo'shiladi.

Chuqurroq muammo — egalik. So'rov bir necha yuz millisekund yashaydi. "Bu odam welcome xatini olishiga ishonch hosil qil" degan ish esa undan uzoqroq yashashi kerak: provayder uzilishidan ham, deploy'dan ham, jarayonning almashtirilishidan ham nariga o'tib. Bunday umrni request handler ichida ifodalab bo'lmaydi. Handler tugaydi va niyatlarini o'zi bilan olib ketadi.

Demak, so'rov javob berishdan oldin rost bo'lishi shart bo'lgan minimumni bajaradi, qolganini esa doimiy yozuv sifatida topshiradi:

@Post('signup')
async signup(@Body() dto: SignupDto) {
  const user = await this.users.create(dto);
 
  // Queue'ga yozish ma'no jihatidan xuddi shu commit chegarasining bir qismi:
  // agar bu throw qilsa, chaqiruvchi xatoni ko'radi va bemalol qayta urina oladi.
  await this.emailQueue.add(
    'welcome',
    { userId: user.id },                       // ID, user obyekti emas
    { jobId: `welcome:${user.id}` },           // dedupe kaliti — quyida tushuntiraman
  );
 
  return { id: user.id };
}

Javob endi tez va halol. Xat baribir ketadi, ertami-kechmi, va "ertami-kechmi" — queue aynan ta'minlash uchun yaratilgan xususiyat.

At-least-once — bu tag, o'chirib qo'yiladigan sozlama emas

Men ishlagan har bir durable queue — BullMQ, SQS, RabbitMQ — at-least-once yetkazib berishni beradi. Mualliflar dangasa bo'lgani uchun emas, balki mustaqil ravishda buziladigan jarayonlar bilan tarmoq ustida exactly-once'ni queue sizga uzatib berolmaydi. Har doim shunday lahza bo'ladi: job bajarilgan, lekin ack yozilmagan. Va o'sha nuqtada faqat ikki variant qoladi — "ehtimol ikki marta bajarish" yoki "ehtimol umuman bajarmaslik".

Queue'lar birinchisini tanlaydi. Bu to'g'ri tanlov va u yukni sizga o'tkazadi: consumer'ingiz idempotent bo'lishi shart. "Imkon bo'lsa yaxshi bo'lardi" emas. Shart. Himoyasiz kartadan pul yechadigan worker ertami-kechmi ikki marta yechadi, chunki worker ertami-kechmi charge bilan ack orasida o'ladi.

Xat holatida himoya — bazadagi da'vo:

async function sendWelcomeOnce(userId: string) {
  // Avval da'vo qil. Insert'da kim yutsa, yuborish o'shaniki.
  const claimed = await db.$executeRaw`
    INSERT INTO email_sends (user_id, kind)
    VALUES (${userId}, 'welcome')
    ON CONFLICT (user_id, kind) DO NOTHING
  `;
 
  if (claimed === 0) return;   // allaqachon yuborilgan yoki yuborilyapti — bu urinish bo'sh o'tadi
 
  await mailer.sendWelcome(userId);
}

Producer'dagi jobId ham foyda beradi, lekin u boshqa muammoni yechadi: ayni o'sha jobning ikki marta navbatga qo'yilishiga yo'l qo'ymaydi. Ayni o'sha job'ning ikki marta yetkazilishiga esa hech qanday ta'siri yo'q — at-least-once kafolati aynan shu holat sodir bo'lishini bildiradi.

Ack qilingan bilan tugallangan bir narsa emas

Eng nozik xatolarning ko'pi shu farq ichida yashaydi.

Worker job'ni olganda, queue uni tugagan deb hisoblamaydi — jarayonda deb biladi va soatni ishga tushiradi. BullMQ'da job waitdan activega o'tadi va uning lock'i har lockRenewTimeda yangilanadi. Faqat processor funksiyangiz resolve bo'lgandagina job completedga ko'chadi. Agar jarayon job active turganda o'lsa, hech kim hech narsani resolve qilmaydi, lock muddati tugaydi va queue job'ni stalled deb qaytarib oladi.

Aynan shu qaytarib olish tufayli o'sha ikki yuz xat printsipial jihatdan tiklanishi mumkin edi. Amalda tiklanmaganining sababi esa: maxStalledCount standart 1 da turgan, lockDuration esa bizning eng sekin yuborishimizdan qisqa edi. Lock'dan uzoqroq davom etadigan va uni yangilamaydigan job "sekin" emas, "o'lgan deb taxmin qilingan" hisoblanadi — va queue uni boshqa birovga beradi, holbuki asl worker uni hamon bemalol bajarayotgan bo'ladi.

const worker = new Worker(
  'email',
  async (job) => {
    await sendWelcomeOnce(job.data.userId);
  },
  {
    connection,
    concurrency: 10,
    // O'rtachadan emas, job tanasining p99 qiymatidan uzunroq.
    lockDuration: 60_000,
    // Stalled job failed'ga tushgunicha necha marta tiklanishi mumkin.
    maxStalledCount: 3,
    stalledInterval: 30_000,
  },
);

Endi men o'ylab ham o'tirmay qo'llaydigan ikkita qoida bor. lockDurationni medianadan emas, job'ning sekin dumidan olib belgila. Va agar job haqiqatan lock'dan uzoq ishlasa, uni job.extendLock() yoki job.updateProgress() bilan ochiq-oydin yangila — sukut o'limdan farq qilmaydi.

Backoff va nega qat'iy intervallar uzilishni battar qiladi

"5 soniyadan keyin qayta urin" degan retry siyosati zararsiz ko'rinadi — nosozlik yuqori oqimda va hamma uchun umumiy bo'lguncha.

Pochta provayderi ikki daqiqaga o'chdi. Besh yuzta job bir soniya ichida yiqildi. Barcha besh yuztasi besh soniyadan keyin, birgalikda qayta uradi. Hammasi birga yiqiladi, besh soniyadan keyin, yana birga uradi. Siz metronom qurdingiz: u qiynalayotgan servisga har besh soniyada sinxron to'lqin urib turadi. Provayder nihoyat tiklanganda esa birinchi ko'radigan narsasi — besh yuzta bir vaqtdagi so'rov, bu esa uni yana yiqitishi mumkin.

Eksponensial backoff urinishlarni vaqt bo'ylab yoyadi. Jitter esa ularni har bir urinish oynasi ichida yoyadi — odamlar aynan shu qismini tashlab ketishadi, holbuki sinxronlikni haqiqatan shu buzadi.

await emailQueue.add('welcome', { userId }, {
  attempts: 5,
  // backoffStrategy'ga faqat 'custom' yo'naltiradi. Ruxsat etilgan
  // qiymatlar: 'fixed' | 'exponential' | 'custom' — bu yerga ixtiyoriy
  // nom yozsangiz, siz o'ylagan narsa jimgina ishlamaydi.
  backoff: { type: 'custom' },
  removeOnComplete: { age: 3600, count: 1000 },
  removeOnFail: false,                  // xatolarni tekshirish uchun saqlab qolamiz
});
 
// Strategiya Queue'da emas, Worker settings'ida turadi.
const worker = new Worker('email', processor, {
  connection,
  settings: {
    backoffStrategy: (attemptsMade: number) => {
      const base = Math.min(1000 * 2 ** attemptsMade, 5 * 60_000);  // 5 daqiqada cheklaymiz
      return Math.floor(base / 2 + Math.random() * (base / 2));     // deyarli to'liq jitter
    },
  },
});

Bu formulada ikki narsa muhim. Chek — chunki cheklanmagan eksponensial backoff oxiri qayta urinishni keyingi haftaga qo'yib yuboradi. Va tasodifiylik — chunki usiz birga yiqilgan besh yuzta job, egri chiziq qanchalik aqlli bo'lmasin, abadiy birga qayta uraveradi.

Retry siyosatining ikkinchi yarmi — nimani qayta urinmaslik kerakligini bilish. Provayderdan kelgan 500 qayta urinishga arziydi. 422 "noto'g'ri email manzil" esa arzimaydi — u besh marta bir xil yiqiladi va siz allaqachon bilgan xulosaga yetish uchun worker quvvatining besh daqiqasini yoqib yuboradi.

try {
  await mailer.sendWelcome(userId);
} catch (err) {
  if (err.status >= 400 && err.status < 500 && err.status !== 429) {
    // Doimiy xato. Qayta urinishni to'xtat va odam ko'rib chiqishi uchun yo'nalt.
    throw new UnrecoverableError(`invalid recipient: ${err.message}`);
  }
  throw err;   // vaqtinchalik — backoff o'z ishini qilsin
}

Dead letter queue — bu bajariladigan ishlar ro'yxati

Urinishlari tugagan job yo'qolib ketmasligi kerak, va hech kim ochib ko'rmagan failed to'plamida yotishi ham kerak emas. Failed to'plami — qaror talab qiladigan ishlar navbati.

Men unga uchta xususiyatga ega operatsion sath sifatida qarayman: u o'sganda alert beradi, har bir yozuvda harakat qilish uchun yetarli kontekst bor va job'ni qaytarib qo'yishning qo'llab-quvvatlanadigan yo'li mavjud.

worker.on('failed', async (job, err) => {
  if (job.attemptsMade < job.opts.attempts) return;   // hali urinilyapti, o'lgani yo'q
 
  logger.error({
    queue: job.queueName,
    jobId: job.id,
    name: job.name,
    data: job.data,          // faqat ID'lar, shuning uchun log qilish xavfsiz
    attempts: job.attemptsMade,
    reason: err.message,
  }, 'job exhausted retries');
 
  metrics.increment('jobs.dead', { queue: job.queueName, name: job.name });
});

Shunda qayta ishga tushirish arxeologik qazishma emas, oddiy skriptga aylanadi:

const dead = await emailQueue.getFailed(0, 500);
for (const job of dead) {
  // attemptsMade standart holatda nolga TUSHMAYDI — `attempts` ni allaqachon
  // tugatgan job wait'ga qaytadi va birinchi xatodayoq yana qulaydi.
  await job.retry('failed', { resetAttemptsMade: true });
}

Replay xavfsiz bo'lishi kerakligining sababi — bu yerdagi qolgan hamma narsanikiga o'xshash: consumer'lar idempotent, shuning uchun qisman bajarilgan job'ni qayta ishga tushirish xavfli emas, zerikarli. Agar DLQ'ni replay qilish qo'rqinchli bo'lsa, muammo DLQ'da emas.

Kichik job'lar, chunki monolitni davom ettirib bo'lmaydi

"40 000 foydalanuvchining oylik hisobotini tayyorla"ni bitta job qilib yozishga kuchli mayl bor. Uni tushunish oson, rejalashtirish oson. Lekin u tiklab bo'lmaydigan: 40 daqiqaning 38-chisida yiqiladi, retry esa nolinchi foydalanuvchidan boshlanadi va har bir retry xuddi o'sha devorga urilish ehtimoli bilan yashaydi.

Muqobili — fan-out. Bitta kichik koordinator job ko'plab kichik birlik job'larni navbatga qo'yadi:

// Koordinator: arzon, tez va ikki marta ishga tushsa ham xavfsiz.
async function scheduleMonthlyReports(job: Job<{ period: string }>) {
  const { period } = job.data;
  let cursor: string | undefined;
 
  do {
    const page = await db.user.findMany({
      where: { active: true, ...(cursor && { id: { gt: cursor } }) },
      orderBy: { id: 'asc' },
      take: 500,
      select: { id: true },
    });
 
    await reportQueue.addBulk(
      page.map((u) => ({
        name: 'report',
        data: { userId: u.id, period },
        opts: { jobId: `report:${period}:${u.id}` },   // har davr uchun har foydalanuvchiga bitta hisobot
      })),
    );
 
    cursor = page.at(-1)?.id;
  } while (cursor);
}

Endi qulash sizga bitta foydalanuvchining hisobotiga tushadi va u bir necha soniyada qayta bajariladi. Jarayon uzoq ishlaydigan tsiklning xotirasida emas, queue'ning o'zida yozilgani uchun barqaror. Parallellikni tekinga olasiz, queue hisoblagichlarini o'qib haqiqiy progress raqamini ko'rasiz, jobId esa koordinatorning o'zi qulaganidan keyin ham uni qayta ishga tushirishni xavfsiz qiladi.

Kelishuv halol: bitta log qatorini 40 000 job'ga aylantirdingiz, bu esa ko'proq Redis xotirasi va ko'proq shovqin. Buni tajovuzkor removeOnComplete hal qiladi. Bu yaxshi kelishuv, lekin baribir kelishuv.

Payload'da ID bo'ladi, obyekt emas

Butun entity'ni payload'ga solish samarali tuyuladi — worker so'rov yubormaydi. Ikki narsa buziladi.

Payload — navbatga qo'yish paytidagi surat. Agar job yigirma daqiqa backlog'da yotsa yoki yiqilib bir soatdan keyin qayta ursa, worker allaqachon mavjud bo'lmagan haqiqat versiyasi ustida ish qiladi. Men foydalanuvchi allaqachon to'g'rilab qo'ygan manzilga xat ketganini ko'rganman — to'g'rilash job navbatga qo'yilgandan keyin bo'lgan ekan.

Yana: payload'lar Redis'da, serializatsiya qilingan holda, har bir job va har bir retry uchun saqlanadi. O'n kilobaytlik ichki obyekt katta queue'ga ko'paytirilsa — bu haqiqiy xotira, va uni aynan eng noqulay paytda sezasiz.

// Mo'rt: yomon eskiradigan va xotira yeydigan surat.
await queue.add('welcome', { user, template, attachments });
 
// Barqaror: identifikatorlar va o'zgarmasligi shart bo'lgan kichik faktlar.
await queue.add('welcome', {
  userId: user.id,
  templateVersion: 'welcome_v4',   // ma'lumotni emas, xatti-harakatni qotiramiz
});

Men ishlatadigan chegara shunday: payload ishni nima aniqlashini va uning ma'nosini nima qotirishini olib yuradi. O'zgaruvchan hamma narsa bajarilish paytida, haqiqat manbaidan o'qiladi.

Ataylab o'chirish

O'sha ikki yuz xatga qaytamiz. Asosiy sabab queue emas edi — bizning jarayonimiz SIGTERMni "hozir o'l" deb tushundi, "yangi ish olishni to'xtat va qo'lingdagini tugat" deb emas.

async function shutdown(signal: string) {
  logger.info({ signal }, 'shutting down');
 
  // close() ning o'zi allaqachon yumshoq yo'l: yangi job olishni to'xtatadi
  // va aktivlarini kutadi. close(true) esa majburlaydi va kutmaydi.
  await worker.close();
  await queue.close();
  await redis.quit();
 
  process.exit(0);
}
 
process.on('SIGTERM', () => void shutdown('SIGTERM'));
process.on('SIGINT', () => void shutdown('SIGINT'));

NestJS'da xuddi shu narsa shutdown hook'larni yoqib, worker'ni modul yopishiga ruxsat berish orqali chiqadi:

@Injectable()
export class EmailWorker implements OnApplicationShutdown {
  async onApplicationShutdown() {
    await this.worker.close();   // jarayon chiqishidan oldin qo'ldagi job'larni tugatadi
  }
}
// main.ts
app.enableShutdownHooks();

Buni haqiqiy qiladigan ikkita shart bor. Orkestratoringizning grace period'i eng uzun job'ingizdan katta bo'lishi kerak — agar Kubernetes 30 soniyadan keyin SIGKILL qilsa, job esa 90 soniya olsa, graceful shutdown xatti-harakat emas, oddiy izoh. Va muddatga chindan ulgurib bo'lmaganda, yuqoridagi stalled job tiklanishiga qaytasiz. Graceful shutdown — tez yo'l; stall recovery — kafolat.

Queue qachon kerak emas

"Handler ichida await sendEmail"ning teskari nosozligi — hamma narsani navbatga tiqish, va uning ham o'z narxi bor: qo'shimcha sakrash, kuzatiladigan qo'shimcha tizim va hech narsa bo'layotganga o'xshamaydigan, keyin esa sirli tarzda bo'lib qoladigan foydalanuvchi tajribasi.

Ish tez bo'lsa va chaqiruvchiga natija chindan kerak bo'lsa — sinxron qoldiring. Hozir qaytaradigan qatorni yozish. Kuponni tekshirish. Javobda URL'i beriladigan bitta kichik rasmni kichraytirish. 30 millisekundlik amalni job'ga o'rash sizga kerak bo'lmagan eventual consistency'ni va tushuntirishingizga to'g'ri keladigan "nega hali paydo bo'lmadi" xatosini beradi.

Ish sekin bo'lsa, siz nazorat qilmaydigan tizim bilan gaplashsa, restart'dan omon chiqishi shart bo'lsa yoki chaqiruvchidan mustaqil qayta urinilishi mumkin bo'lsa — navbatga qo'ying. Agar javob chindan "foydalanuvchi buning muvaffaqiyati yoki muvaffaqiyatsizligini hozir ko'rishi shart" bo'lsa, queue noto'g'ri vosita va frontend'dagi hech qancha polling uni to'g'riga aylantirmaydi.

Bu aslida nimani mashq qildiradi

Bu yerdagi har bir usul bitta harakatning turli ko'rinishi: jarayondagi ishning holatini barqaror va tashqaridan ko'rinadigan qil, toki jarayonning o'limi o'chirib tashlash emas, uzilish bo'lsin.

Idempotent consumer'lar, ochiq ack, stall recovery, jitter bilan backoff, kichik va davom ettiriladigan birliklar, faqat ID'li payload'lar, ataylab o'chirish — bularning birortasi ham qulashning oldini olmaydi. Ular qulashni qiziqarsiz qiladi. Worker istalgan lahzada o'lishga haqli, tizimning javobi esa qaysidir qatordagi biroz kechroq timestamp bo'lishi kerak, support ticket emas.

Men endi har qanday background pipeline'ga qo'yadigan sinov — bitta savol: hozir worker'ni -9 bilan o'ldirsam, tizim nimani yo'qotadi? Agar javob "bir necha soniya" bo'lsa — ish tugagan. Agar javobni bilish uchun kodni o'qish kerak bo'lsa — tugamagan.

O'xshash maqolalar

backendpayments

Soya rejimi: kodni ishonishdan oldin isbotlash

Joytop'dagi yangi to'lov tasdiqlash kodi bir hafta productionda soya rejimida ishladi — qaror qabul qilmasdan, faqat kuzatib. Ertaga u yagona qaror qiluvchiga aylanadi. Soya rejimi nima, u nimadan himoya qiladi va nega uni kuzatuvsiz yoqib qo'yib bo'lmaydi.

3 daqiqa o'qish
backendredis

Keshni bekor qilish — amaliyotda

Hazil bu ikkita qiyin muammodan biri deydi. Haqiqat esa shuki, ko'p jamoalar bekor qilish strategiyasini umuman yozmaydi — ular TTL yozib, umid qiladi. Muqobillar aslida nimaga tushishini ko'ramiz.

8 daqiqa o'qish