GĐ10 — Queues, Jobs, Workers & Cronjob: xử lý nền đúng cách

GĐ09 đã giới thiệu BullMQ và outbox ở mức "chạy được". Giai đoạn này đi sâu vào những thứ làm sập production lúc 2 giờ sáng: job chạy hai lần, cron chạy trên cả 5 instance, worker bị kill giữa chừng, queue tắc nghẽn vì một tenant, và job "đã xong" nhưng dữ liệu chỉ hoàn thành một nửa.

Đây là giai đoạn phân biệt rõ nhất giữa "biết dùng thư viện" và "hiểu hệ thống bất đồng bộ". Nó cũng là nền tảng bắt buộc trước khi vào GĐ19 — Distributed Systems.

Kiểm chứng ngày 2026-10-05: BullMQ 6.x (chạy thử 6.3.11, cài kèm ioredis); repeat đã bị xoá ở v6, dùng upsertJobScheduler; Redis cho queue cần maxmemory-policy noeviction theo docs going to production; limiter là của cả queue, rate limit theo nhóm là BullMQ Pro (docs). Redis 8.6.1 và PostgreSQL 17.9 là bản đã dùng để chạy thử các ví dụ có nhãn "đã chạy".


1. Vì sao cần xử lý nền#

Vấn đề. Request HTTP có ngân sách thời gian: người dùng chờ, load balancer timeout (thường 30–60s), trình duyệt bỏ cuộc. Nhưng nhiều việc không vừa ngân sách đó: gửi 10.000 email, resize video, gọi LLM 40 giây, sinh báo cáo PDF.

Giải pháp. Tách "nhận yêu cầu" khỏi "làm việc":

textReady
Client ──POST──► API ──đẩy job──► Queue ──lấy job──► Worker ──► kết quả         ◄─202──┘  (trả về ngay)                    (process riêng)

Ba lợi ích, nói cho rõ:

  1. Độ trễ — người dùng nhận 202 Accepted trong 20ms thay vì chờ 40 giây.
  2. Chống chịu — worker chết thì job vẫn nằm trong queue, chạy lại được. Nếu làm inline, request chết là mất việc.
  3. Điều tiết tải — spike 10.000 request không giết DB; queue giữ lại và worker rút với tốc độ bạn kiểm soát.

Cái giá. Bạn đổi một hệ thống đồng bộ dễ hiểu lấy một hệ thống bất đồng bộ khó debug: job chạy hai lần, chạy sai thứ tự, chạy 3 tiếng sau, hoặc không chạy mà không ai biết. Phần còn lại của giai đoạn này là về việc trả cái giá đó cho đúng.

Việc nào nên đẩy nền:

NênKhông nên
Gửi email/SMS/pushXác thực, phân quyền
Xử lý ảnh/video, sinh PDFBất cứ thứ gì người dùng cần thấy kết quả ngay
Gọi API bên thứ ba chậmGhi dữ liệu mà bước sau phụ thuộc trực tiếp
Tính toán tổng hợp, báo cáoViệc nhỏ hơn ~50ms (chi phí queue lớn hơn lợi ích)
Đồng bộ sang search index

2. Chọn công nghệ hàng đợi#

Công nghệMô hìnhĐảm bảoDùng khi
BullMQ (Redis)Job queueAt-least-onceNode monolith. Mặc định của lộ trình này
pg-boss (Postgres)Job queue trên SQLAt-least-onceĐã có Postgres, không muốn thêm Redis
Bảng SQL tự viếtFOR UPDATE SKIP LOCKEDAt-least-onceTải nhỏ, muốn hiểu cơ chế, muốn job trong cùng transaction với dữ liệu
RabbitMQMessage broker (AMQP)At-least-once, routing mạnhNhiều service, cần routing/fanout phức tạp
KafkaLog phân tánAt-least-once, giữ thứ tự trong partition, replay đượcEvent streaming, throughput rất lớn, nhiều consumer group
SQS / Cloud TasksManagedAt-least-once (FIFO queue: exactly-once trong 5 phút)Không muốn vận hành hạ tầng

Khác biệt quan trọng nhất: queue vs log.

  • Queue (BullMQ, SQS, Rabbit): job được tiêu thụ rồi biến mất. Một job → một worker xử lý.
  • Log (Kafka): sự kiện được giữ lại, nhiều consumer group đọc độc lập, mỗi group có offset riêng, replay được từ đầu.

Nếu bạn cần "phát lại 3 ngày sự kiện vì service mới vừa lên", bạn cần Kafka chứ không phải BullMQ. Nếu bạn cần "gửi email này một lần", queue đủ và đơn giản hơn nhiều.

Pitfall #1 — chọn Kafka cho DA3. Kafka giải quyết vấn đề của hệ thống có nhiều đội và throughput rất lớn. Với một monolith, nó là chi phí vận hành lớn không đổi lấy gì. Chọn BullMQ, và biết vì sao bạn không chọn Kafka — đó mới là câu trả lời phỏng vấn tốt.

Redis cho BullMQ: một instance riêng, noeviction. Đừng dùng chung Redis với cache (GĐ09 mục 10). Cache cần tự xoá khi đầy bộ nhớ (allkeys-lru); queue cần điều ngược lại: BullMQ yêu cầu maxmemory-policy noeviction (docs BullMQ), vì nếu Redis âm thầm xoá khoá của queue thì job và trạng thái của nó biến mất. Với noeviction, khi đầy bộ nhớ Redis từ chối ghi (lỗi rõ ràng) thay vì xoá dữ liệu. Bật thêm AOF cho instance này (cũng theo docs). BullMQ 6 không kéo sẵn ioredis: cài npm i bullmq ioredis (đã chạy thật trên 6.3.11: thiếu ioredis thì new Queue báo lỗi).


3. Đảm bảo giao nhận: at-most-once / at-least-once / exactly-once#

Đây là khái niệm quan trọng nhất của cả giai đoạn.

Đảm bảoNghĩaĐánh đổi
At-most-onceChạy 0 hoặc 1 lầnCó thể mất việc. Ack trước khi làm
At-least-onceChạy 1 hoặc nhiều lầnCó thể lặp việc. Ack sau khi làm
Exactly-onceĐúng 1 lầnKhông tồn tại ở tầng giao nhận qua mạng

Vì sao exactly-once là không thể. Worker xử lý xong, gửi ACK, rồi chết trước khi ACK tới nơi. Queue không phân biệt được "chết trước khi làm" với "chết sau khi làm, trước khi báo". Nó buộc phải chọn: giao lại (→ lặp) hoặc bỏ qua (→ mất). Không có lựa chọn thứ ba.

Kết luận thực hành: mọi hệ thống thật đều là at-least-once, và bạn đạt hiệu quả exactly-once bằng cách làm handler idempotent. Đây là câu trả lời phỏng vấn chuẩn, và cũng là điều bạn phải thực sự làm trong code.

Sơ đồ: ACK mất và lựa chọn của queue
textReady
 Worker                         Queue/Broker   | --- nhận job ------------------> |   | (làm việc: gửi mail, charge)     |   | --- ACK ---------X (mất / chết)  |   Broker chỉ thấy: không có ACK   |                                  |   |   Hai khả năng mà broker không phân biệt được:   |     (1) worker chết TRƯỚC khi làm   -> cần giao lại   |     (2) worker chết SAU khi làm     -> giao lại sẽ làm LẦN HAI   |   Chọn giao lại  => at-least-once (có thể lặp)   <- mọi hệ thống thật   Chọn bỏ qua    => at-most-once  (có thể mất)

Hệ quả: giao lại là mặc định an toàn, và việc "không làm hai lần" chuyển sang handler (idempotency). Bảng kết quả mô phỏng ở ngay dưới là bằng chứng cho từng mẫu.

Ba cách làm handler idempotent:

A. Máy trạng thái PENDING → DONE ở tầng DB. Cách "ghi khoá processed_job trước, rồi mới gọi Stripe" trông chắc chắn nhưng thực ra là at-most-once: nếu lỗi tạm xảy ra sau khi ghi khoá và trước khi charge xong, lần retry gặp khoá có sẵn, tưởng "đã làm rồi" và thoát êm. Giao dịch không bao giờ xảy ra. Khoá phải có hai trạng thái và có hạn:

textReady
CREATE TABLE idempotency_key (  key          TEXT PRIMARY KEY,  status       TEXT NOT NULL CHECK (status IN ('PENDING', 'DONE')),  locked_until TIMESTAMPTZ NOT NULL,        -- PENDING chỉ giữ chỗ đến thời điểm này  lock_token   UUID NOT NULL,               -- mỗi lần giành khoá sinh token mới; chỉ người giữ token mới nhả/ghi được  result       JSONB);
typescriptReady
// Stripe client đặt timeout/retry tường minh: mặc định stripe-node là timeout 80s và 2 lần retry mạngconst stripe = new Stripe(process.env.STRIPE_KEY!, { timeout: 15_000, maxNetworkRetries: 2 })// lockSeconds ≥ timeout × (1 + maxNetworkRetries) + backoff giữa các lần retry + đệm//            = 15s × 3 + (≤ 2 × 5s) + đệm ⇒ 75s. Muốn khoá ngắn hơn: đặt maxNetworkRetries: 0 (lúc đó 15s + đệm).const LOCK_SECONDS = 75// Giành khoá: trả lock_token nếu ta được làm; null = đã DONE hoặc người khác đang giữasync function claim(key: string, lockSeconds: number): Promise<string | null> {  const rows = await db.$queryRaw<{ lock_token: string }[]>`    INSERT INTO idempotency_key (key, status, locked_until, lock_token)    VALUES (${key}, 'PENDING', now() + make_interval(secs => ${lockSeconds}), gen_random_uuid())    ON CONFLICT (key) DO UPDATE      SET locked_until = now() + make_interval(secs => ${lockSeconds}), lock_token = gen_random_uuid()      WHERE idempotency_key.status = 'PENDING' AND idempotency_key.locked_until < now()    RETURNING lock_token`  return rows[0]?.lock_token ?? null}async function handleChargeJob(job: Job<ChargePayload>) {  const key = `charge:${job.data.orderId}`  const token = await claim(key, LOCK_SECONDS)  if (!token) {    const row = await db.idempotencyKey.findUnique({ where: { key } })    if (row?.status === 'DONE') return                        // đã xong → thoát êm    throw new Error('khoá đang được giữ')                     // PENDING còn hạn → để BullMQ retry sau  }  let intent  try {    // Cùng khoá gửi cho provider: lớp bảo vệ cho khe hở bên dưới    intent = await stripe.paymentIntents.create(params, { idempotencyKey: key })  } catch (e) {    // Lỗi bắt được → nhả khoá NGAY để retry (backoff 1s, 2s, 4s...) không phải chờ hết LOCK_SECONDS.    // An toàn vì provider khử trùng theo key: dù lỗi là timeout sau khi Stripe đã nhận, retry không trừ tiền lần hai.    // Điều kiện `lock_token`: nếu khoá của ta đã hết hạn và worker khác giành lại, câu này ảnh hưởng 0 dòng    // thay vì mở khoá của họ cho người thứ ba.    await db.$executeRaw`      UPDATE idempotency_key SET locked_until = now()      WHERE key = ${key} AND lock_token = ${token}::uuid AND status = 'PENDING'`    throw e  }  // 0 dòng = khoá đã bị giành lại: người giữ mới sẽ gọi Stripe cùng key, nhận cùng kết quả rồi tự ghi DONE  await db.$executeRaw`    UPDATE idempotency_key SET status = 'DONE', result = ${JSON.stringify({ id: intent.id })}::jsonb    WHERE key = ${key} AND lock_token = ${token}::uuid AND status = 'PENDING'`}

Hành vi (đã chạy: luồng claim → charge → DONE bằng nhà cung cấp thanh toán giả lập khử trùng theo key, các câu SQL trên PostgreSQL 17.9 thật với nhiều connection):

Kịch bảnGhi khoá rồi mới chargePENDING → DONE
Lỗi tạm (provider 503) sau khi giữ khoá, rồi retry ngayretry bị bỏ qua, 0 lần trừ tiềnbắt lỗi, nhả khoá; retry ngay: charge 1 lần; retry nữa: DONE, bỏ qua
Timeout sau khi provider đã nhận tiền, rồi retry ngay—nhả khoá; retry gọi lại cùng key, provider khử trùng: 1 lần trừ tiền
20 worker cùng khoá—1 charge, 19 chờ, 1 lần trừ tiền
Khoá cũ: A hết hạn, B giành, A nhả khoá bằng token cũ, C claimA nhả xoá luôn khoá của B, C giành được: hai người cùng giữA nhả/ghi DONE bằng token cũ: 0 dòng; C bị từ chối; B ghi DONE: 1 dòng
Câu claim, nhả khoá và DONE trên PostgreSQL 17 thật, nhiều connection—20 claim song song: đúng 1 thắng; 5 worker chạy đủ luồng: 1 charge; lỗi tạm rồi retry: 1 charge; sau DONE (kể cả quá hạn khoá): claim 0

Giới hạn của cách A — nói thẳng.

  • Nếu process chết giữa stripe.paymentIntents.create và UPDATE ... DONE, đến khi khoá hết hạn worker khác sẽ charge lần hai (mô phỏng: 2 lần trừ tiền). DB không biết được charge đã xảy ra hay chưa. Vì vậy cách A luôn đi kèm cách C: cùng key gửi sang provider (header Idempotency-Key của Stripe, ≤255 ký tự, Stripe giữ ít nhất 24 giờ, cùng key nhưng khác tham số thì lỗi). Cùng mô phỏng, khi provider khử trùng theo key: 1 lần trừ tiền. Tính đúng đắn = trạng thái trong DB cộng khoá phía provider.
  • Điều kiện của "1 lần trừ tiền": nó chỉ đúng trong thời hạn Stripe còn giữ key (tối thiểu 24 giờ) và khi tham số gọi y hệt lần đầu. Job nằm trong DLQ rồi được chạy lại sau thời hạn đó, đúng lúc process từng chết giữa charge và DONE, thì Stripe coi là key mới và có thể trừ tiền lần hai. Vì vậy replay DLQ cho job thanh toán nên đối chiếu PaymentIntent bên Stripe (tra theo orderId trong metadata) trước khi chạy lại, hoặc chạy lại trong cửa sổ 24 giờ.
  • locked_until phải lớn hơn thời gian xử lý xấu nhất: timeout × (1 + số lần retry mạng của client) + backoff + đệm (code trên: 75 giây). Ngắn hơn thì worker thứ hai giành lại khoá trong khi worker đầu còn đang charge; mặc định stripe-node (80s × 3 lần thử) đòi khoá hơn 4 phút nếu không chỉnh client. Khi khoá vẫn hết hạn lúc worker cũ còn chạy (process bị treo, GC dài), lock_token giữ cho worker cũ không nhả hay ghi DONE lên khoá của người mới (chỉ chặn việc ghi nhầm, không thay thế việc chọn lockSeconds đủ lớn: hai worker vẫn có thể cùng gọi Stripe, và key phía provider mới là thứ khử trùng).
  • Nếu không nhả khoá khi lỗi, retry gặp "khoá đang giữ" suốt lockSeconds. Với attempts: 5 và backoff 1/2/4/8 giây (tổng ~15 giây, mục 5) mà khoá giữ 90 giây, cả bốn retry đều bị từ chối, job vào failed và tiền không bao giờ bị trừ: chính lỗi "mất giao dịch" mà cách A sinh ra để tránh. Vì vậy handler nhả khoá khi bắt được lỗi (code trên).
  • Process chết thì không ai nhả khoá: các retry trong khoảng khoá còn hạn vẫn gặp "đang giữ". Chọn lockSeconds ≥ timeout × (1 + maxNetworkRetries) + backoff + đệm (như tính ở trên), và hoặc để tổng backoff phủ qua khoảng đó, hoặc chấp nhận job vào DLQ rồi chạy lại có chủ đích (an toàn nhờ khoá phía provider).
  • Các câu SQL (make_interval, ON CONFLICT ... WHERE ... RETURNING lock_token, gen_random_uuid()) đã chạy trên PostgreSQL 17.9 thật, nhiều connection tranh chấp; chưa chạy qua Prisma $queryRaw/$executeRaw (phần ::uuid, ::jsonb trong code là để tham số khớp kiểu cột, chưa kiểm), chưa gọi Stripe thật (nhà cung cấp là bản giả lập). Mặc định timeout 80s và 2 lần retry mạng của stripe-node đọc từ mã nguồn bản 23.0.0.
typescriptReady
// B. Thao tác tự nhiên idempotent — tốt nhất khi làm đượcawait db.user.update({ where: { id }, data: { status: 'ACTIVE' } })   // gán, không phải $incawait s3.putObject({ Key: deterministicKey, Body: data })             // ghi đè cùng key// C. Đẩy idempotency sang bên thứ baawait stripe.paymentIntents.create(params, { idempotencyKey: `order-${orderId}` })

Pitfall #2 — dùng jobId làm khoá idempotency. BullMQ tự dọn job đã hoàn thành (removeOnComplete), nên jobId có thể được cấp lại và "đã xử lý chưa" không còn tra được. Khoá idempotency phải sống trong DB của bạn, với vòng đời bạn kiểm soát.

Pitfall #3 — $inc trong job. UPDATE counters SET n = n + 1 chạy hai lần cho kết quả sai. Nếu bắt buộc phải cộng dồn, ghi sự kiện có id duy nhất rồi tổng hợp, đừng cộng trực tiếp.


4. Transactional Outbox — bài toán hai-lần-ghi#

Vấn đề. Bạn cần làm hai việc nguyên tử ở hai hệ thống khác nhau:

typescriptReady
await db.$transaction(async (tx) => {  await tx.order.create({ data: order })  await queue.add('send-confirmation', { orderId })   // ← SAI, cả hai chiều})

Sai theo chiều thứ nhất: queue.add() thành công, rồi transaction rollback → job gửi email cho một đơn hàng không tồn tại. Sai theo chiều thứ hai: transaction commit, rồi Redis chết trước khi queue.add() xong → đơn hàng tồn tại nhưng không ai được báo, vĩnh viễn.

Giải pháp: outbox. Ghi "ý định gửi" vào cùng database, cùng transaction. Một tiến trình riêng đọc bảng đó và đẩy sang queue.

textReady
CREATE TABLE outbox (  id           BIGSERIAL PRIMARY KEY,  topic        TEXT        NOT NULL,  payload      JSONB       NOT NULL,  created_at   TIMESTAMPTZ NOT NULL DEFAULT now(),  published_at TIMESTAMPTZ,  attempts     INT         NOT NULL DEFAULT 0);CREATE INDEX idx_outbox_pending ON outbox (id) WHERE published_at IS NULL;
typescriptReady
// 1. Ghi cùng transaction — nguyên tử thật sự vì cùng một DBawait db.$transaction(async (tx) => {  const order = await tx.order.create({ data: input })  await tx.outbox.create({    data: { topic: 'order.created', payload: { orderId: order.id } },  })})// 2. Relay chạy riêng — SKIP LOCKED cho phép nhiều instance chạy song song an toànasync function relayOnce() {  await db.$transaction(async (tx) => {    const rows = await tx.$queryRaw<OutboxRow[]>`      SELECT * FROM outbox      WHERE published_at IS NULL      ORDER BY id      LIMIT 100      FOR UPDATE SKIP LOCKED`    for (const row of rows) {      await queue.add(row.topic, row.payload, {        jobId: `outbox-${row.id}`,          // chống lặp ở tầng queue        removeOnComplete: 1000,      })      await tx.$executeRaw`UPDATE outbox SET published_at = now() WHERE id = ${row.id}`    }  }, { timeout: 30_000 })   // Prisma 7 mặc định 5 giây; khoá được giữ suốt lúc gọi Redis}

FOR UPDATE SKIP LOCKED giải quyết gì. Nó cho phép relay #2 bỏ qua các hàng relay #1 đang giữ, thay vì xếp hàng chờ. Đây là nền tảng của mọi queue xây trên SQL. Không có nó, chạy nhiều relay là vô nghĩa.

Outbox cho bạn at-least-once, không phải exactly-once — relay có thể chết sau queue.add trước UPDATE. Vì vậy handler vẫn phải idempotent (mục 3). Hai cơ chế này bổ sung nhau, không thay thế nhau.

Sơ đồ và kết quả mong đợi: relay chết giữa chừng
textReady
 API (1 transaction)           Relay (loop)                 BullMQ        Worker INSERT order  ──┐             SELECT ... FOR UPDATE INSERT outbox ──┴─COMMIT ───► SKIP LOCKED  (khoá hàng)                               queue.add(jobId=outbox-7) ─► job ─► handler                               UPDATE published_at              (idempotent)                               COMMIT Chết sau queue.add, trước COMMIT => ROLLBACK, hàng vẫn "chưa publish" => relay kế tiếp add LẠI cùng jobId => BullMQ bỏ qua bản trùng    (chỉ khi job cũ còn tồn tại)

Đã chạy (BullMQ 6.3.11, Redis 8.6.1, PostgreSQL 17.9): 5 hàng trong outbox; relay thêm 3 job rồi ném lỗi trước COMMIT → published = 0, queue có 3 job. Chạy lại relay → published = 5, queue có đúng 5 job (3 job trùng jobId không được thêm lần hai). Cẩn thận: chống trùng bằng jobId chỉ hiệu lực khi job cũ còn trong Redis. Nếu job đã removeOnComplete xong rồi relay mới replay hàng cũ, job được thêm lại (đã chạy: removeOnComplete: true, add lại cùng jobId sau khi job đầu xong → handler chạy 2 lần). Handler idempotent (mục 3) mới là chốt chặn cuối.

Khi nào dùng outbox: khi mất sự kiện gây hậu quả thật (thanh toán, đơn hàng, email giao dịch, đồng bộ trạng thái). Khi nào không cần: analytics, log không quan trọng, việc có thể tính lại được.

Dọn dẹp. Bảng outbox mọc mãi. Xoá hàng published_at < now() - interval '7 days' bằng cron, hoặc partition theo ngày (→ GĐ14).


5. Thiết kế job: payload, retry, timeout#

Payload phải nhỏ và ổn định.

typescriptReady
// TỐT — id + phiên bảnawait queue.add('resize-image', { assetId: '01H...', v: 1 })// XẤU — nhét cả object; dữ liệu đã cũ khi worker chạy, và có thể chứa PIIawait queue.add('resize-image', { asset: { ...30 fields... }, user: { email, phone } })

Ba lý do: (1) payload là ảnh chụp lúc enqueue, có thể lỗi thời khi worker chạy 3 phút sau; (2) payload nằm trong Redis, PII trong đó là rủi ro tuân thủ; (3) payload lớn tốn RAM Redis và làm chậm mọi thứ.

Thêm v (version) vào payload ngay từ đầu. Khi bạn đổi hình dạng payload, các job cũ vẫn còn trong queue với hình dạng cũ. Có v thì handler xử lý được cả hai; không có thì bạn phải xả queue hoặc chấp nhận job lỗi.

Retry với exponential backoff + jitter:

typescriptReady
await queue.add('send-email', payload, {  attempts: 5,  backoff: { type: 'exponential', delay: 1000 },  // 1s, 2s, 4s, 8s, 16s  removeOnComplete: { age: 3600, count: 1000 },  removeOnFail:     { age: 86400 },})

Vì sao cần jitter. Nếu 5.000 job cùng fail vì một API sập rồi cùng retry sau đúng 1 giây, bạn tạo ra một đợt tấn công vào chính API vừa hồi phục — gọi là thundering herd. Jitter (ngẫu nhiên hoá độ trễ) làm chúng dàn ra. BullMQ 6 có sẵn tuỳ chọn jitter (0 đến 1) trong backoff, ví dụ backoff: { type: 'exponential', delay: 1000, jitter: 0.5 }; cần công thức riêng thì dùng backoff strategy tuỳ biến.

Phân biệt lỗi tạm thời và lỗi vĩnh viễn. Đây là thứ hầu hết code bỏ qua:

typescriptReady
try {  await sendEmail(job.data)} catch (e) {  if (isPermanent(e)) {            // 400 Bad Request, email không hợp lệ, record bị xoá    await markFailedPermanently(job.data)    return                          // KHÔNG throw → không retry vô ích  }  throw e                           // 429, 5xx, timeout mạng → để BullMQ retry}

Retry một địa chỉ email sai cú pháp 5 lần là lãng phí và làm nhiễu cảnh báo.

Kết quả mong đợi: jitter và lỗi vĩnh viễn

Với BullMQ 6 có sẵn backoff: { type: 'exponential', delay, jitter } (jitter từ 0 đến 1, ngẫu nhiên hoá phần độ trễ). Để job dừng ngay thay vì return êm, ném UnrecoverableError của BullMQ (job vào failed, không retry; xem ví dụ ở bài tập mục 11).

typescriptReady
const opts = { attempts: 4, backoff: { type: 'exponential', delay: 300, jitter: 0.5 } }

Đã chạy (BullMQ 6.3.11): job lỗi tạm (503) chạy 4 lần, khoảng cách giữa các lần [277, 506, 1118] ms (gốc 300/600/1200 ms; mỗi lần rơi ngẫu nhiên trong khoảng [gốc × (1 − jitter), gốc], theo mã nguồn 6.3.11, nên không lần nào cố định). Job lỗi vĩnh viễn chạy đúng 1 lần rồi vào failed. Sai thường gặp: nghĩ jitter làm khoảng chờ dài hơn gốc; nó chỉ rút ngắn.

Timeout. Mọi lời gọi ra ngoài trong job phải có timeout riêng: một lời gọi treo mãi giữ một slot concurrency mãi, và mọi slot treo thì worker đứng yên. Timeout này không cần ngắn hơn lockDuration: BullMQ tự gia hạn khoá khi worker còn sống (xem Pitfall #4).

typescriptReady
new Worker('emails', handler, {  connection,  concurrency: 10,  // Mặc định (BullMQ 6, đọc từ mã nguồn 6.3.11): lockDuration 30_000,  // lockRenewTime = lockDuration / 2, stalledInterval 30_000, maxStalledCount 1})// handler tự đặt timeout cho lời gọi ngoàiconst res = await fetch(url, { signal: AbortSignal.timeout(30_000) })

Khoá job hoạt động thế nào. Khi worker lấy job, nó giữ một khoá trên job đó (thời hạn lockDuration) và tự gia hạn mỗi lockRenewTime (mặc định một nửa lockDuration). Job chạy 5 phút vẫn giữ được khoá miễn là worker còn gia hạn được. Cứ mỗi stalledInterval, các worker kiểm tra những job active có khoá đã hết hạn; job như vậy gọi là stalled và:

  • bị đưa về hàng đợi (wait) để worker khác chạy lại từ đầu, stalledCounter tăng 1;
  • nếu vượt maxStalledCount thì chuyển sang failed với lý do job stalled more than allowable limit.

Đã chạy thật trên BullMQ 6.3.11 (Redis 8, lockDuration 2000, stalledInterval 1000):

  • job chạy 7 giây bằng setTimeout (không chặn event loop): chạy đúng 1 lần, lockRenewTime tự đặt 1000;
  • worker A chặn event loop 6 giây bằng vòng lặp đồng bộ: khoá hết hạn, worker B (process khác) nhận job lúc ~3,9 giây với stalledCounter = 1 — job chạy hai lần; khi A thoát vòng lặp và báo hoàn tất, nó nhận lỗi Missing lock for job / Lock mismatch, nhưng tác dụng phụ của lần chạy A đã xảy ra;
  • khi cả A lẫn B đều bị chặn, lần stalled thứ hai làm job failed với job stalled more than allowable limit.

Pitfall #4 — job chạy hai lần vì stalled. Đây là nguyên nhân số một của "job của tôi chạy hai lần mà tôi không hiểu vì sao". Nguyên nhân không phải "job lâu hơn lockDuration" (khoá được gia hạn tự động), mà là worker không gia hạn được:

  • Event loop bị chặn quá lockDuration: vòng lặp CPU đồng bộ dài (nén, băm, parse JSON khổng lồ).
  • Worker chết hoặc mất kết nối Redis (bị SIGKILL, OOM, GC pause rất dài).

Cách xử lý: (1) handler idempotent (mục 3), vì stalled chạy lại là bình thường và tác dụng phụ của lần chạy cũ không rút lại được; (2) việc nặng CPU không chạy thẳng trong handler: dùng sandboxed processor (process/thread riêng) hoặc chia nhỏ để trả quyền cho event loop; (3) nếu bị chặn ngắn hợp lệ, tăng lockDuration cho rộng hơn chỗ chặn. job.updateProgress() chỉ ghi tiến độ, không phải cơ chế gia hạn khoá (đọc mã nguồn 6.3.11).

5.1 Job trễ, flow cha-con và limiter#

Ba công cụ của BullMQ cho những việc mục 5 chưa cover. Cả ba đều sống trong Redis, nên cùng chịu luật noeviction ở mục 2.

typescriptReady
import { FlowProducer, Queue, Worker } from 'bullmq'// Job trễ: chạy sau 24 giờ, không phải lúc chính xác đến từng giâyawait reminders.add('remind', { orderId: 'o1' }, { delay: 24 * 3600_000 })// Flow: 2 job con ở queue 'thumb', xong cả hai thì job cha chạy ở queue 'report'const flows = new FlowProducer({ connection })await flows.add({  name: 'build-report', queueName: 'report', data: { reportId: 'r1' },  children: [    { name: 'thumb', queueName: 'thumb', data: { id: 'a' }, opts: { failParentOnFailure: true } },    { name: 'thumb', queueName: 'thumb', data: { id: 'b' }, opts: { failParentOnFailure: true } },  ],})new Worker('report', async (job) => {  const values = await job.getChildrenValues()   // { 'bull:thumb:<id>': kết quả từng job con }  return { got: Object.values(values) }}, { connection })// Limiter: tối đa 10 job mỗi giây cho CẢ queue, dù có bao nhiêu workernew Worker('emails', handler, { connection, limiter: { max: 10, duration: 1000 } })// Thủ công, khi API ngoài trả 429: dừng cả queue 30 giây, job không bị tính là lỗinew Worker('call-api', async () => {  const res = await fetch('https://api.example.com/x')  if (res.status === 429) {    await apiQueue.rateLimit(30_000)    throw Worker.RateLimitError()  }}, { connection, limiter: { max: 1000, duration: 1000 } })   // cần có limiter.max để cơ chế bật

Đã chạy (BullMQ 6.3.11, Redis 8.6.1; code trên qua tsc --strict sạch, các số dưới đo trên bản chạy tương đương):

ThửKết quả
delay: 1500 cho một job, thêm cùng lúc một job không trễjob không trễ chạy ngay; job trễ chạy sau khoảng 1,6 giây
Flow, hai con thành côngcha chạy một lần, getChildrenValues() có đủ kết quả hai con
Con thứ nhất ném UnrecoverableError, không đặt tuỳ chọncha kẹt ở trạng thái waiting-children (đã quan sát sau 1,5 giây)
Cùng ca, failParentOnFailure: truecha vào failed, waitUntilFinished bị reject
Cùng ca, ignoreDependencyOnFailure: truecha vẫn chạy, getChildrenValues() chỉ có kết quả của con thành công
limiter: { max: 3, duration: 1000 }, hai worker, 9 job3, 3, 3 job mỗi giây: giới hạn tính cho cả queue
queue.rateLimit(1500) rồi throw Worker.RateLimitError()job đó chạy lại sau khoảng 1,5 giây, failed bằng 0

Điều cần nhớ:

  • Flow không phải saga. Cha chạy khi mọi con xong; không có bước bồi hoàn. Mỗi job con vẫn phải idempotent (mục 3). Hai tuỳ chọn còn lại, removeDependencyOnFailure và continueParentOnFailure, có trong kiểu của 6.3.11 (chưa chạy). Theo docs, cha và con có thể ở queue khác nhau, jobId không được chứa dấu :, và xoá cha thì xoá luôn các con (docs flows).
  • delay không phải lịch. Docs nói job không bảo đảm chạy đúng thời điểm, còn tuỳ worker đang bận (docs delayed). Việc lặp lại theo lịch dùng scheduler ở mục 8. Job đã delayed sửa thời gian được bằng job.changeDelay (theo docs, chưa chạy).
  • limiter là của cả queue, không phải của từng tenant. Một tenant ồn ào vẫn chiếm hết hạn mức của queue. Rate limit theo nhóm là tính năng BullMQ Pro (group: { limit: { max, duration } } ở worker, docs Pro; docs OSS ghi group key đã bị bỏ từ 3.0: docs rate limiting); với bản OSS, vẫn dùng cách ở mục 9. worker.rateLimit() bị đánh dấu deprecated trong kiểu của 6.3.11: dùng queue.rateLimit().

6. Dead Letter Queue và quan sát#

DLQ là nơi job hết số lần retry đi đến. Nguyên tắc: job không bao giờ được biến mất im lặng.

typescriptReady
worker.on('failed', async (job, err) => {  // Hết số lần retry, HOẶC lỗi vĩnh viễn (UnrecoverableError dừng ngay ở lần 1, attemptsMade < attempts)  // name: đề phòng lỗi qua ranh giới process (sandboxed processor) hoặc hai bản bullmq, nơi instanceof sai  const permanent = err instanceof UnrecoverableError || err.name === 'UnrecoverableError'  const final = job && (job.attemptsMade >= (job.opts.attempts ?? 1) || permanent)  if (job && final) {    await deadLetter.add('failed-job', {      queue: job.queueName, name: job.name, data: job.data,      error: err.message, stack: err.stack, failedAt: new Date().toISOString(),    })    logger.error({ jobId: job.id, queue: job.queueName, err: err.message }, 'job dead-lettered')    metrics.increment('job.dead_letter', { queue: job.queueName })  }})

Đã chạy (BullMQ 6.3.11, Redis 8.6.1): job perm ném UnrecoverableError rồi job temp hết 3 lần retry. Điều kiện cũ (attemptsMade >= attempts một mình) chỉ đưa temp vào DLQ; perm thất bại ở lần 1 nên bị bỏ sót im lặng. Điều kiện trên đưa cả hai vào DLQ (dlqSize: 2). UnrecoverableError có name === 'UnrecoverableError' (đã chạy 6.3.11); so khớp thêm theo name là phòng xa cho lỗi đi qua ranh giới process (suy luận, chưa thử).

Bốn chỉ số phải theo dõi — và ngưỡng cảnh báo:

Chỉ sốNghĩaCảnh báo khi
Queue depth (waiting)Số job đang chờTăng đơn điệu trong 15 phút → worker không theo kịp
Job age (oldest waiting)Job cũ nhất chờ bao lâuVượt SLA nghiệp vụ (vd: email > 5 phút)
Failure rateTỉ lệ fail / tổng> 1% hoặc tăng đột biến
DLQ sizeSố job chết> 0 là phải có người xem

Queue depth là chỉ số quan trọng nhất và cũng hay bị bỏ quên nhất. Nó là tín hiệu sớm: throughput worker < throughput enqueue. Nếu chỉ nhìn "job có fail không", bạn sẽ phát hiện vấn đề khi email đã trễ 6 tiếng.

Bảng điều khiển. bull-board cho BullMQ — nhưng phải đặt sau xác thực. Nó cho phép xem payload (có thể có PII) và chạy lại job.

Chạy lại từ DLQ. Phải là hành động có chủ đích của con người, không tự động. Trước khi chạy lại: đã sửa nguyên nhân chưa? Handler có idempotent không? Dữ liệu còn hợp lệ không (job 3 ngày tuổi có thể tham chiếu record đã xoá)?


7. Worker: vòng đời và graceful shutdown#

Worker phải là process riêng, không phải cùng process với API. Lý do:

  • Job nặng CPU chặn event loop → API treo (→ GĐ03).
  • Scale độc lập: 2 API + 8 worker, hoặc ngược lại.
  • Deploy độc lập; worker crash không làm sập API.
typescriptReady
// worker.ts — entrypoint riêng, Dockerfile chung, command khácconst worker = new Worker('emails', handler, { connection, concurrency: 10 })let shuttingDown = falseasync function shutdown(signal: string) {  if (shuttingDown) return                     // SIGTERM + SIGINT có thể đến cùng lúc  shuttingDown = true  logger.info({ signal }, 'shutting down worker')  // Thoát cưỡng bức: phải NHỎ HƠN grace period của orchestrator (đã trừ preStop)  setTimeout(() => process.exit(1), 20_000).unref()  await worker.close()          // ngừng nhận job mới, CHỜ job đang chạy xong  await db.$disconnect()        // đóng dependency SAU khi worker đã drain  await connection.quit()  process.exit(0)}process.on('SIGTERM', () => void shutdown('SIGTERM'))process.on('SIGINT',  () => void shutdown('SIGINT'))

Đây là cùng mẫu graceful shutdown ở GĐ09 mục 18 (đủ sáu quy tắc ở đó), chỉ thay server.close() bằng worker.close(). Đã chạy thật trên BullMQ 6.3.11: với một job 1,5 giây đang chạy, await worker.close() trả về sau ~1,1 giây và job hoàn tất. Nếu hẹn giờ cưỡng bức nổ trước khi job xong, process chết giữa chừng: khoá job hết hạn, job thành stalled và chạy lại ở worker khác (mục 5), một lý do nữa để handler luôn idempotent.

Sơ đồ và kết quả mong đợi: worker nhận SIGTERM
textReady
 SIGTERM ──► shuttingDown = true ──► hẹn giờ thoát cưỡng bức (20 s, unref)                  │                  ▼        worker.close()  ── ngừng nhận job mới, CHỜ job đang chạy xong                  │                  ▼        db.$disconnect() ──► connection.quit() ──► process.exit(0) Hẹn giờ nổ trước khi job xong ──► exit(1) ──► khoá job hết hạn ──► job stalled                                                ──► worker khác chạy lại từ đầu

Đã chạy (Node 24.21, BullMQ 6.3.11): worker có job 1,5 giây, gửi SIGTERM sau 0,8 giây. Log theo thứ tự got SIGTERM, job done, close() sau 821 ms, exit code 0. Dấu hiệu sai: thấy exit=1 hoặc không có dòng job done (job bị cắt).

Vì sao SIGTERM quan trọng. Khi bạn deploy, orchestrator (Docker, ECS, K8s) gửi SIGTERM rồi chờ một khoảng grace period trước khi SIGKILL. Nếu worker không xử lý SIGTERM, job đang chạy bị cắt giữa chừng → dữ liệu dở dang.

Hai con số phải khớp nhau, và một con số không liên quan:

textReady
thời gian job dài nhất  <  grace period của orchestrator (trừ thời gian preStop)

Nếu grace period (K8s mặc định 30s) ngắn hơn job dài nhất, mọi lần deploy đều cắt job; khoá của nó hết hạn sau tối đa cỡ lockDuration + stalledInterval (suy luận từ cơ chế ở mục 5; mặc định khoảng một phút) rồi job chạy lại từ đầu. Chỉnh terminationGracePeriodSeconds → GĐ18. lockDuration không nằm trong bất đẳng thức này: nó không phải trần thời gian của job (khoá được gia hạn tự động), chỉ quyết định bao lâu sau khi worker chết thì job được giao lại.

concurrency là bao nhiêu? Không có con số ma thuật. Bắt đầu ở 5–10, rồi đo:

  • Job I/O-bound (gọi API, chờ DB): concurrency cao được, giới hạn thật là connection pool DB.
  • Job CPU-bound: concurrency > số core là vô ích, còn làm chậm đi.

Pitfall #5 — concurrency vượt connection pool. 20 worker × concurrency 10 = 200 kết nối cùng lúc tới Postgres. Pool 10 → 190 job xếp hàng chờ connection và timeout. Tính tổng: số worker × concurrency ≤ pool size (hoặc dùng PgBouncer → GĐ21).

7.1 Ghi chú vận hành BullMQ trên production#

Nguồn: docs going to production (đối chiếu ngày 2026-10-05), cộng hai điều đã chạy trên 6.3.11. Redis noeviction và AOF đã ở mục 2.

typescriptReady
// API (producer): thất bại nhanh khi Redis chết thay vì treo requestconst producerConn = new Redis(process.env.QUEUE_REDIS_URL!, { enableOfflineQueue: false })const emails = new Queue('emails', {  connection: producerConn,  defaultJobOptions: {                       // mặc định BullMQ giữ job xong/lỗi mãi mãi    removeOnComplete: { age: 3600, count: 1000 },    removeOnFail: { age: 7 * 86400 },  },})emails.on('error', (err) => console.error({ err: err.message }, 'queue error'))// Worker: maxRetriesPerRequest phải là null; giữ hàng đợi offline mặc địnhconst workerConn = new Redis(process.env.QUEUE_REDIS_URL!, { maxRetriesPerRequest: null })const worker = new Worker('emails', handler, { connection: workerConn })worker.on('error', (err) => console.error({ err: err.message }, 'worker error'))process.on('unhandledRejection', (reason) => { console.error({ reason }, 'unhandledRejection') })process.on('uncaughtException', (err) => { console.error({ err }, 'uncaughtException'); process.exit(1) })
ViệcVì saoBằng chứng
Đặt removeOnComplete/removeOnFail (nên đặt một lần ở defaultJobOptions)Mặc định job xong và job lỗi nằm lại mãi; trên Redis noeviction, bộ nhớ đầy thì Redis từ chối ghi, kể cả queue.add (suy luận từ noeviction)docs; đã chạy: sau một job xong và một job lỗi, getJobCounts còn completed: 1, failed: 1
enableOfflineQueue: false cho producer, giữ mặc định cho workerRedis chết giữa chừng: producer báo lỗi ngay để API trả 503 hoặc để outbox (mục 4) giữ việc, thay vì treodocs; đã chạy: queue.add bị reject sau 0 ms với Stream isn't writeable and enableOfflineQueue options is false; mặc định thì treo ít nhất 3 giây
maxRetriesPerRequest: null cho kết nối workerWorker chờ Redis quay lại thay vì ném lỗidocs
Gắn on('error') cho Worker và QueueLỗi kết nối Redis phải vào log, không im lặngdocs
Bắt unhandledRejection và uncaughtExceptionLỗi không bắt được làm worker chết giữa job (job sẽ stalled, mục 5)docs; process.exit(1) sau uncaughtException là lựa chọn của bài này
Payload ở dạng bản rõ (plaintext) trong RedisMã hoá dữ liệu nhạy cảm trước khi enqueue, hoặc chỉ đưa id (mục 5)docs

Code trên qua tsc --strict sạch. Giới hạn của enableOfflineQueue: false: nó chỉ làm lệnh bị từ chối nhanh khi Redis chết sau khi đã kết nối. Nếu Redis không chạy ngay từ lúc khởi động, queue.add không bị reject mà cũng không hoàn tất (đã thử với cổng đóng: chưa settle sau 5 giây), nên giữ một timeout riêng quanh add hoặc kiểm /health/ready bằng redis.ping().


8. Cronjob và scheduled job#

Đây là phần yếu nhất trong hầu hết codebase. Mọi người viết setInterval(cleanup, 3600_000) rồi scale lên 3 instance, và cleanup chạy 3 lần.

Bốn cách chạy việc định kỳ:

CáchƯuNhược
setInterval trong appĐơn giản nhấtChạy trên mọi instance; mất khi restart; trôi thời gian
BullMQ Job Scheduler (upsertJobScheduler)Chỉ một instance chạy (Redis điều phối); có retry, có lịch sửPhụ thuộc Redis; cần hiểu key ổn định của scheduler
Cron hệ thống / K8s CronJobTách hẳn khỏi app; đúng một lầnCần hạ tầng riêng; khó quan sát trong app
Managed scheduler (EventBridge, Cloud Scheduler)Không vận hành gìKhoá vào nhà cung cấp; gọi qua HTTP nên phải bảo vệ endpoint

Khuyến nghị cho DA3: BullMQ Job Scheduler — cùng hạ tầng, cùng retry, cùng observability với job thường.

typescriptReady
await queue.upsertJobScheduler(  'nightly-cleanup',                       // key ổn định → không sinh lịch trùng  { pattern: '0 3 * * *', tz: 'Asia/Ho_Chi_Minh' },  { name: 'cleanup', data: {} },)

Pitfall #6 — lịch trùng lặp. Với API repeatable cũ (đã bị xoá ở BullMQ 6: queue.add(..., { repeat }) không còn tạo lịch, job chỉ chạy một lần ngay; đã chạy thật trên 6.3.11), mỗi lần app khởi động lại mà repeatJobKey khác đi là bạn có thêm một lịch nữa. Sau 10 lần deploy, job chạy 10 lần mỗi đêm. Luôn dùng key ổn định (upsertJobScheduler), và liệt kê lịch hiện có lúc boot để kiểm tra:

typescriptReady
const schedulers = await queue.getJobSchedulers()logger.info({ count: schedulers.length, keys: schedulers.map(s => s.key) }, 'cron schedulers')

Kết quả mong đợi (đã chạy, BullMQ 6.3.11): gọi upsertJobScheduler('nightly-cleanup', ...) ba lần liên tiếp rồi getJobSchedulers() trả về đúng một khoá: [ 'nightly-cleanup' ]. Nếu boot ra count tăng sau mỗi lần deploy, khoá của bạn không ổn định.

Múi giờ. Luôn khai báo tz tường minh. Server chạy UTC, người dùng ở Asia/Ho_Chi_Minh; "báo cáo hàng ngày lúc 3 giờ sáng" là 3 giờ sáng của ai? Và cron 0 2 * * * ở múi giờ có DST có thể bị bỏ hoặc chạy hai lần vào ngày chuyển giờ. Chi tiết → GĐ14 mục về timezone.

8.1 Leader election — chạy đúng một lần trên N instance#

Khi không dùng được BullMQ (ví dụ việc phải chạy trong process API), cần khoá phân tán:

typescriptReady
// Khoá Redis đơn giản — đủ dùng cho việc không quan trọng sống chếtasync function withLock(key: string, ttlMs: number, fn: () => Promise<void>) {  const token = crypto.randomUUID()  const ok = await redis.set(key, token, 'PX', ttlMs, 'NX')   // NX = chỉ đặt nếu chưa có  if (!ok) return                                             // instance khác đang giữ  try {    await fn()  } finally {    // Chỉ xoá nếu token vẫn là của mình — tránh xoá nhầm khoá của instance khác    await redis.eval(      `if redis.call("get", KEYS[1]) == ARGV[1] then return redis.call("del", KEYS[1]) else return 0 end`,      1, key, token,    )  }}await withLock('cron:cleanup', 5 * 60_000, cleanup)

Cảnh báo trung thực về khoá Redis. Cơ chế trên không an toàn tuyệt đối: nếu tiến trình giữ khoá bị treo (GC pause, network partition) quá TTL, instance khác lấy được khoá trong khi tiến trình đầu vẫn tưởng mình đang giữ → hai bên cùng chạy. Đây là vấn đề nổi tiếng (tranh luận Redlock giữa Martin Kleppmann và Salvatore Sanfilippo, 2016) và không có cách sửa hoàn toàn ở tầng khoá.

Cách xử lý đúng: đừng dựa vào khoá để đảm bảo tính đúng đắn. Dùng khoá để giảm việc trùng (tối ưu hoá), và dựa vào idempotency để đảm bảo đúng đắn. Nếu thực sự cần đúng-một-lần cứng, dùng hệ thống có consensus thật: etcd, ZooKeeper, hoặc lease của Kubernetes (→ GĐ19 mục về consensus).

Cách đơn giản nhất và thường là đúng nhất: dùng khoá ở tầng DB, đúng ngữ nghĩa:

textReady
-- Advisory lock của Postgres: tự nhả khi session kết thúc, kể cả khi process chếtSELECT pg_try_advisory_lock(hashtext('cron:cleanup'));

Cạnh sắc: khoá cấp session đi qua connection pool. "Session kết thúc" nghĩa là connection đóng, mà connection trong pool sống rất lâu. Đã chạy (PostgreSQL 17.9 tạm, driver pg, pool 3): giành khoá ở connection A rồi pg_advisory_unlock ở connection B trả false kèm cảnh báo you don't own a lock of type ExclusiveLock; A trả về pool mà chưa nhả thì khoá vẫn còn (một connection khác thử lấy: false) cho đến khi connection đó bị đóng. Với pool (Prisma, pg.Pool), dùng khoá cấp transaction, tự nhả khi COMMIT, ROLLBACK hoặc mất connection:

typescriptReady
import type { Pool, PoolClient } from 'pg'async function runExclusive(pool: Pool, key: string, fn: (tx: PoolClient) => Promise<void>) {  const tx = await pool.connect()  try {    await tx.query('BEGIN')    const { rows } = await tx.query<{ ok: boolean }>(      'SELECT pg_try_advisory_xact_lock(hashtext($1)) AS ok', [key])    if (!rows[0].ok) { await tx.query('ROLLBACK'); return false }   // instance khác đang chạy    await fn(tx)    await tx.query('COMMIT')                                        // khoá nhả ở đây    return true  } catch (e) {    await tx.query('ROLLBACK').catch(() => {})    throw e  } finally {    tx.release()  }}

Kết quả mong đợi (đã chạy, ba instance gọi cùng lúc, việc 300 ms): một đã chạy, hai bỏ qua; gọi lại sau COMMIT thì chạy được. tsc --strict sạch. Ba lưu ý: (1) transaction giữ một connection suốt lúc chạy, nên chỉ hợp việc ngắn; với Prisma $transaction interactive mặc định 5 giây (xem mục 4), việc dài hơn phải nới timeout hoặc dùng BullMQ scheduler thay vì khoá (chưa chạy qua Prisma); (2) hashtext trả số 32 bit và là hàm nội bộ (chưa xác minh nó có trong tài liệu PostgreSQL hay được bảo đảm giữa các bản): hai khoá khác tên trùng số chỉ khiến hai việc không liên quan chạy tuần tự, không sai dữ liệu; (3) giống khoá Redis, khoá này giảm việc trùng, không thay idempotency.

8.2 Missed run và catch-up#

Server sập từ 2:00 đến 4:00, cron 3:00 không chạy. Chuyện gì xảy ra khi server lên?

Ba chính sách, chọn có chủ đích:

  1. Bỏ qua — hợp với việc idempotent chạy thường xuyên (dọn cache mỗi giờ). Mặc định của BullMQ.
  2. Chạy bù một lần — hợp với việc tổng hợp (báo cáo ngày).
  3. Chạy bù tất cả — hầu như luôn sai; sập 2 ngày là 48 lần chạy dồn, gây quá tải ngay lúc vừa hồi phục.

Thiết kế tốt hơn cron thuần: job theo khoảng dữ liệu, không theo thời điểm. Thay vì "chạy lúc 3h sáng để xử lý hôm qua", hãy lưu con trỏ:

typescriptReady
// Job biết mình đã xử lý tới đâu → tự bù, chạy lại vô hại, không phụ thuộc lịchconst { lastProcessedAt } = await db.jobCursor.findUniqueOrThrow({ where: { name: 'daily-rollup' } })const until = startOfDay(new Date())for (const day of eachDayBetween(lastProcessedAt, until)) {  await rollupDay(day)                                        // idempotent  await db.jobCursor.update({ where: { name: 'daily-rollup' }, data: { lastProcessedAt: day } })}

Thiết kế này tự động xử lý missed run, chạy lại được bao nhiêu lần cũng được, và test được mà không cần giả lập thời gian.

Kết quả mong đợi (đã chạy bản tham chiếu ở bài tập mục 11, Vitest 5): con trỏ ở ngày 1/10, server sập đến sáng 5/10 → lần chạy kế tiếp xử lý ngày 2, 3, 4 (mỗi ngày đúng một lần); chạy lại ngay thì xử lý 0 ngày. Lưu ý: nếu process chết giữa rollupDay và update con trỏ, ngày đó chạy lại, nên rollupDay phải ghi đè kết quả (upsert theo ngày), không cộng dồn.


9. Fairness: không để một tenant chiếm cả queue#

Vấn đề "noisy neighbor" phiên bản queue. Tenant A upload 50.000 file. Queue có 50.000 job của A. Tenant B upload 1 file — nằm ở vị trí 50.001 và chờ 4 tiếng.

Ba cách xử lý, từ đơn giản tới phức tạp:

  1. Queue riêng theo mức ưu tiên: emails:critical (OTP, reset password) và emails:bulk (newsletter). Worker riêng cho mỗi queue. Đơn giản, hiệu quả, làm trước tiên.

  2. priority của BullMQ: số nhỏ chạy trước. Nhưng cẩn thận — priority queue có thể gây starvation: job priority thấp không bao giờ tới lượt nếu job priority cao liên tục đổ vào.

  3. Rate limit theo nhóm — đúng nhất cho multi-tenant:

typescriptReady
const worker = new Worker('emails', handler, { connection, concurrency: 20 })// Mỗi tenant tối đa 10 job / phút, phần dư tự động đẩy sang sauawait queue.add('send', payload, { group: { id: `tenant:${tenantId}` } })// (BullMQ Pro có group rate limit; bản OSS: tự chia queue theo tenant hoặc//  dùng token bucket trong Redis kiểm tra ở đầu handler rồi re-enqueue có delay)

Cách thủ công dùng được với BullMQ OSS:

typescriptReady
async function handler(job: Job) {  const allowed = await tokenBucket.take(`tenant:${job.data.tenantId}`, 10, 60_000)  if (!allowed) {    await queue.add(job.name, job.data, { delay: 10_000 })   // trả lại queue, thử sau    return  }  await doWork(job.data)}

Kết quả mong đợi (suy luận từ code, chưa chạy): tenant A đẩy 50.000 job; mỗi phút chỉ 10 job của A chạy, phần còn lại quay lại queue với delay: 10_000, nên job của tenant B không còn chờ sau 50.000 job. Cái giá: job bị bỏ lại tốn thêm một vòng qua Redis, và return sau queue.add đánh dấu job gốc completed (đếm metric cho đúng).


10. Test job và worker#

Nguyên tắc: tách handler khỏi hạ tầng queue. Handler là một hàm thuần nhận payload — test nó trực tiếp, không cần Redis.

typescriptReady
// handler là hàm độc lập, không biết BullMQ tồn tạiexport async function handleSendEmail(payload: SendEmailPayload, deps: Deps) { ... }// worker.ts chỉ nối dâynew Worker('emails', (job) => handleSendEmail(job.data, deps), { connection })// test — không Redis, không workerit('không gửi lại email đã gửi', async () => {  await handleSendEmail(payload, deps)  await handleSendEmail(payload, deps)          // chạy lần hai  expect(deps.mailer.send).toHaveBeenCalledTimes(1)   // idempotent})

Test tích hợp queue (Redis thật qua Testcontainers → GĐ13): kiểm tra job được đẩy vào queue với payload đúng, kiểm tra retry, kiểm tra đường vào DLQ.

Không dùng sleep cố định để chờ job xong. Dùng job.waitUntilFinished(queueEvents, ttl): promise trả kết quả của job, bị reject khi job fail hoặc hết ttl.

typescriptReady
const events = new QueueEvents('emails', { connection })     // một lần cho cả file testbeforeAll(() => events.waitUntilReady())                     // chưa sẵn sàng thì bỏ lỡ sự kiệnafterAll(() => events.close())it('job gửi mail hoàn tất', async () => {  const job = await queue.add('send', payload)  expect(await job.waitUntilFinished(events, 5_000)).toEqual({ sent: payload.to })})

Đã chạy (Vitest 5, BullMQ 6.3.11, Redis 8.6.1, tsc --strict sạch): job thường trả đúng kết quả của handler; job ném UnrecoverableError('bad input') làm waitUntilFinished bị reject với bad input; job còn delayed quá ttl bị reject với "timed out before finishing". Mẫu cũ new QueueEvents(...).on('completed', resolve) hết hạn cả hai ca đã thử: job xong trước khi QueueEvents kịp kết nối, và add job ngay sau khi tạo listener mà không chờ waitUntilReady(). Nó còn resolve theo bất kỳ job nào hoàn tất và không bao giờ đóng QueueEvents.

Test cron bằng cách gọi thẳng hàm, không chờ lịch. Với logic phụ thuộc thời gian, tiêm clock thay vì gọi new Date() trực tiếp — cùng nguyên tắc như test timezone ở GĐ14.


11. Bài tập — nâng cấp DA3#

Yêu cầu.

  1. Tách worker thành process riêng, Dockerfile chung nhưng command khác; docker-compose chạy api + worker + redis + postgres.

    Lời giải và cách kiểm tra

    Hướng làm: một Dockerfile (stage build + stage chạy node:24-alpine), hai service dùng cùng image, khác command. Redis cho queue tách khỏi Redis cache (mục 2). Code tham chiếu (chưa chạy: không dùng Docker của người dùng):

    textReady
    services:  postgres:    image: postgres:18    environment: { POSTGRES_PASSWORD: dev }    healthcheck: { test: ["CMD-SHELL", "pg_isready -U postgres"], interval: 3s }  redis-queue:    image: redis:8    command: ["redis-server", "--maxmemory-policy", "noeviction", "--appendonly", "yes"]    healthcheck: { test: ["CMD", "redis-cli", "ping"], interval: 3s }  api:    build: .    command: ["node", "dist/server.js"]    depends_on: { postgres: { condition: service_healthy }, redis-queue: { condition: service_healthy } }  worker:    build: .    command: ["node", "dist/worker.js"]    stop_grace_period: 30s    depends_on: { postgres: { condition: service_healthy }, redis-queue: { condition: service_healthy } }

    Kết quả mong đợi: docker compose up -d --build rồi docker compose ps thấy 4 service running; docker compose stop worker thì curl localhost:3000/health/live của API vẫn 200 (worker chết không kéo API chết). Lỗi hay gặp: command không đổi nên worker chạy cả API; Redis queue không noeviction; API khởi động trước khi Postgres sẵn sàng.

  2. Ba queue: critical (OTP, reset password), default (email giao dịch), bulk (newsletter, export) — worker và concurrency riêng cho mỗi queue.

    Lời giải và cách kiểm tra

    Hướng làm: mỗi queue một Worker với concurrency riêng; tổng worker × concurrency không vượt pool DB (mục 7, Pitfall #5). critical cao để OTP không chờ, bulk thấp để không chiếm DB. Code tham chiếu (đã chạy: tsc --noEmit; chưa chạy riêng đo độ trễ):

    typescriptReady
    import { Queue, Worker, type Processor } from 'bullmq'import { Redis } from 'ioredis'export const connection = new Redis(process.env.QUEUE_REDIS_URL!, { maxRetriesPerRequest: null })const plan = { critical: 20, default: 10, bulk: 2 } as constexport const queues = Object.fromEntries(  Object.keys(plan).map((n) => [n, new Queue(n, { connection })]),)export function startWorkers(handlers: Record<keyof typeof plan, Processor>) {  return (Object.keys(plan) as (keyof typeof plan)[]).map(    (n) => new Worker(n, handlers[n], { connection, concurrency: plan[n] }),  )}

    Kết quả mong đợi (suy luận): đẩy 5.000 job bulk rồi 1 job critical, job critical chạy ngay vì worker của nó rảnh. Lỗi hay gặp: maxRetriesPerRequest: null bị quên khi tự tạo ioredis (BullMQ yêu cầu cho Worker); dùng chung một Worker cho cả ba queue.

  3. Outbox cho mọi sự kiện quan trọng; relay chạy trong worker process với SKIP LOCKED; cron dọn outbox cũ hơn 7 ngày.

    Lời giải và cách kiểm tra

    Hướng làm: ghi outbox cùng transaction với dữ liệu (mục 4). Relay là một việc lặp trong worker process: dùng upsertJobScheduler('outbox-relay', { every: 1000 }, ...) để chỉ một instance chạy mỗi nhịp, hoặc một vòng while có sleep (không setInterval chồng nhịp). Cron dọn: DELETE FROM outbox WHERE published_at < now() - interval '7 days'. Sơ đồ và code relay: mục 4 (hình "relay chết giữa chừng"). Kết quả mong đợi (đã chạy, PostgreSQL 17.9): relay chết trước COMMIT → published = 0 và 3 job trong queue; relay chạy lại → published = 5, đúng 5 job. Câu xoá nhận 3 hàng (8 ngày tuổi, 1 ngày tuổi, chưa publish) chỉ xoá đúng hàng 8 ngày: trả 1, còn lại 2. Lỗi hay gặp: xoá cả hàng chưa publish (đừng thêm OR published_at IS NULL vào câu xoá: hàng chưa publish là việc chưa giao); jobId không đủ để chống trùng khi removeOnComplete đã xoá job cũ (mục 4).

  4. Handler idempotent bằng bảng idempotency_key (PENDING → DONE, mục 3) cộng khoá idempotency gửi sang provider; chứng minh bằng test gọi handler hai lần, và một lần với lỗi tạm sau khi giành khoá.

    Lời giải và cách kiểm tra

    Hướng làm: dùng nguyên mẫu PENDING → DONE ở mục 3, tách hai thứ qua tham số để test được: handleChargeJob(pool, charge, orderId), trong đó charge(key) gọi provider với idempotencyKey: key. Test dùng provider giả khử trùng theo key (giống Stripe) và PostgreSQL thật cho bảng idempotency_key. Code tham chiếu (đã chạy, Vitest 5, PostgreSQL 17.9, kết quả 4 passed):

    typescriptReady
    // provider giả khử trùng theo keyfunction makeProvider() {  const seen = new Map<string, { id: string }>()  let charges = 0, failNext = 0  return {    failNext(n: number) { failNext = n },    get charges() { return charges },    async charge(key: string) {      if (failNext > 0) { failNext--; throw Object.assign(new Error('503'), { status: 503 }) }      const hit = seen.get(key)      if (hit) return hit      charges++      const r = { id: `pi_${charges}` }      seen.set(key, r)      return r    },  }}it('gọi hai lần: một lần trừ tiền', async () => {  const p = makeProvider()  expect(await handleChargeJob(pool, (k) => p.charge(k), 'o1')).toBe('charged')  expect(await handleChargeJob(pool, (k) => p.charge(k), 'o1')).toBe('skipped')  expect(p.charges).toBe(1)})it('lỗi tạm sau khi giành khoá: retry ngay vẫn trừ đúng một lần', async () => {  const p = makeProvider()  p.failNext(1)  await expect(handleChargeJob(pool, (k) => p.charge(k), 'o2')).rejects.toThrow('503')  expect(p.charges).toBe(0)                    // chưa trừ, nhưng khoá đã nhả  expect(await handleChargeJob(pool, (k) => p.charge(k), 'o2')).toBe('charged')  expect(await handleChargeJob(pool, (k) => p.charge(k), 'o2')).toBe('skipped')  expect(p.charges).toBe(1)})it('20 worker cùng khoá: một charge, 19 busy', async () => {  const p = makeProvider()  const slow = async (k: string) => { await new Promise((r) => setTimeout(r, 100)); return p.charge(k) }  const res = await Promise.allSettled(Array.from({ length: 20 }, () => handleChargeJob(pool, slow, 'o3')))  const ok = res.filter((r) => r.status === 'fulfilled').length  const busy = res.filter((r) => r.status === 'rejected' && r.reason instanceof LockBusyError).length  expect([ok, busy, p.charges]).toEqual([1, 19, 1])})

    Kết quả mong đợi: npx vitest run → Tests 4 passed (4) (bản chạy có thêm một ca: token cũ không nhả được khoá của người giữ mới). Bài test lỗi tạm phải đỏ trên mẫu cũ (ghi khoá processed_job rồi mới charge): đã thử trên PostgreSQL thật, mẫu cũ cho p.charges = 0 và lần retry trả skipped, nên dòng expect(p.charges).toBe(1) sẽ thất bại. Lỗi hay gặp: mock luôn trả thành công nên không bao giờ chạm nhánh lỗi tạm; chạy các ca song song mà dùng chung key; test provider không khử trùng nên "1 lần trừ tiền" không còn đúng (xem giới hạn ở mục 3: chết giữa charge và DONE).

  5. Retry có exponential backoff + jitter; phân biệt lỗi tạm thời và vĩnh viễn.

    Lời giải và cách kiểm tra

    Hướng làm: attempts + backoff.type: 'exponential' + jitter; phân loại lỗi một chỗ, lỗi vĩnh viễn thì ném UnrecoverableError để không retry. Code tham chiếu (đã chạy):

    typescriptReady
    import { UnrecoverableError } from 'bullmq'export class PermanentError extends Error {}export function toBullError(e: unknown): unknown {  if (e instanceof PermanentError) return new UnrecoverableError(e.message)  const status = (e as { status?: number }).status  if (status && status >= 400 && status < 500 && status !== 429 && status !== 408) {    return new UnrecoverableError(`HTTP ${status}`)   // 4xx: gửi lại cũng vậy  }  return e                                             // 429, 408, 5xx, timeout: retry}// handler: try { ... } catch (e) { throw toBullError(e) }await queue.add('send-email', payload, {  attempts: 4,  backoff: { type: 'exponential', delay: 1000, jitter: 0.5 },})

    Kết quả mong đợi (đã chạy với delay: 300): lỗi 503 chạy 4 lần, khoảng cách [277, 506, 1118] ms; PermanentError chạy 1 lần rồi failed. Lỗi hay gặp: coi mọi 4xx là vĩnh viễn (429 và 408 là tạm thời); nuốt lỗi bằng return rồi job hiện completed dù thất bại; không có jitter nên retry đồng loạt (mục 5).

  6. DLQ + log có cấu trúc + endpoint /admin/dlq (sau xác thực) để xem và chạy lại.

    Lời giải và cách kiểm tra

    Hướng làm: sự kiện failed đẩy vào queue dead-letter khi hết retry hoặc lỗi vĩnh viễn (sửa ở mục 6). Endpoint chỉ cho admin: liệt kê, và chạy lại bằng cách add job gốc vào đúng queue sau khi người vận hành đã xem. Code tham chiếu (listener đã chạy; endpoint chưa chạy: Express không cài trong môi trường thử, đây là khung):

    typescriptReady
    worker.on('failed', async (job, err) => {  const permanent = err instanceof UnrecoverableError || err.name === 'UnrecoverableError'  const final = job && (job.attemptsMade >= (job.opts.attempts ?? 1) || permanent)  if (job && final) {    await deadLetter.add('failed-job', {      queue: job.queueName, name: job.name, data: job.data, error: err.message,    })    logger.error({ jobId: job.id, queue: job.queueName, err: err.message }, 'job dead-lettered')  }})router.get('/admin/dlq', requireAdmin, async (_req, res) => {  const jobs = await deadLetter.getJobs(['waiting'], 0, 49)  res.json(jobs.map((j) => ({ id: j.id, ...j.data })))        // payload có thể có PII: chỉ admin})router.post('/admin/dlq/:id/retry', requireAdmin, async (req, res) => {  const j = await deadLetter.getJob(req.params.id)  if (!j) return res.sendStatus(404)  await queues[j.data.queue].add(j.data.name, j.data.data)    // có chủ đích, sau khi đã sửa nguyên nhân  await j.remove()  res.sendStatus(202)})

    Kết quả mong đợi (đã chạy phần listener, BullMQ 6.3.11): job temp (3 lần retry) và job perm (UnrecoverableError) đều vào DLQ, dlqSize: 2. Với điều kiện cũ chỉ có temp. Lỗi hay gặp: DLQ không có xác thực; chạy lại tự động; chạy lại job thanh toán quá 24 giờ mà không đối chiếu provider (mục 3).

  7. Graceful shutdown xử lý SIGTERM theo GĐ09 mục 18, chờ job đang chạy xong; ghi vào README job dài nhất, grace period (trừ preStop) và các giá trị lockDuration/stalledInterval đang dùng.

    Lời giải và cách kiểm tra

    Hướng làm: dùng khung ở mục 7 (xem sơ đồ ở đó), không process.exit trước worker.close(). Kiểm tra bằng kill -TERM lúc job đang chạy. Kết quả mong đợi (đã chạy): got SIGTERM, job done, close() sau 821 ms, exit=0. README cần ghi (số mẫu, thay bằng số thật của bạn): job dài nhất 15 giây; grace period 30 giây trừ preStop 5 giây = 25 giây khả dụng; hẹn giờ thoát cưỡng bức 20 giây (nhỏ hơn phần khả dụng); lockDuration 30000, stalledInterval 30000, maxStalledCount 1 (mặc định BullMQ 6, đọc từ mã nguồn 6.3.11). Lỗi hay gặp: hẹn giờ thoát lớn hơn grace period (bị SIGKILL trước khi tự thoát); đóng DB trước worker.close() nên job đang chạy mất kết nối.

  8. Cron: dọn dẹp hàng đêm (BullMQ scheduler, tz tường minh) + một job rollup theo con trỏ như mục 8.2.

    Lời giải và cách kiểm tra

    Hướng làm: upsertJobScheduler với khoá ổn định và tz tường minh (mục 8); rollup đọc con trỏ, xử lý từng ngày đã kết thúc, ghi đè theo ngày, rồi mới dịch con trỏ. Code tham chiếu (đã chạy, Vitest 5, 3 passed):

    typescriptReady
    const DAY = 86_400_000export const startOfDay = (d: Date) => new Date(Math.floor(d.getTime() / DAY) * DAY)  // UTCexport async function runRollup(  store: { get(): Promise<Date>; set(day: Date): Promise<void> },  rollupDay: (day: Date) => Promise<void>,  now: Date,) {  const until = startOfDay(now)                 // ngày hôm nay chưa kết thúc: không động tới  let processed = 0  for (let t = (await store.get()).getTime() + DAY; t < until.getTime(); t += DAY) {    await rollupDay(new Date(t))                // phải ghi đè, không cộng dồn    await store.set(new Date(t))    processed++  }  return processed}

    Kết quả mong đợi: con trỏ 1/10, now 5/10 03:00 → xử lý 3 ngày (2, 3, 4), mỗi ngày 1 lần; chạy lại → 0; chết giữa ngày 3 rồi chạy lại → ngày 3 chạy 2 lần nhưng con trỏ cuối là 4/10. Đã chạy: upsertJobScheduler ba lần cùng khoá → một lịch ([ 'nightly-cleanup' ]). Lỗi hay gặp: xử lý cả ngày hôm nay (dữ liệu chưa đủ); lẫn "ngày" theo UTC với "ngày" theo múi giờ nghiệp vụ (đổi startOfDay cho khớp múi giờ, → GĐ14); dịch con trỏ trước khi rollupDay xong (mất ngày khi chết giữa chừng).

  9. Metric: queue depth, oldest job age, failure rate, DLQ size — xuất ra /metrics hoặc log; đặt một cảnh báo thật.

    Lời giải và cách kiểm tra

    Hướng làm: gom bốn chỉ số mỗi vài giây, xuất ra /metrics (Prometheus) hoặc log JSON. Tỉ lệ fail lấy từ hai bộ đếm completed/failed tăng dần (do QueueEvents hoặc worker.on), để hệ thống metric tự tính tốc độ. Code tham chiếu (đã chạy hai phép đo getJobCounts và tuổi job cũ nhất):

    typescriptReady
    async function collect(queue: Queue, dlq: Queue) {  const c = await queue.getJobCounts('waiting', 'active', 'failed')  const [oldest] = await queue.getWaiting(0, 0)               // job chờ lâu nhất  return {    depth: c.waiting,    oldestAgeMs: oldest ? Date.now() - oldest.timestamp : 0,    dlqSize: await dlq.count(),  }}

    Kết quả mong đợi (đã chạy): 5 job chờ, không worker, sau ~1,2 giây → { depth: 5, oldestAgeMs: 1209 }. Cảnh báo thật (khung, chưa chạy: không có Prometheus trong môi trường thử): dlq_size > 0 trong 5 phút, hoặc oldest_age_seconds vượt SLA. Lỗi hay gặp: chỉ theo dõi tỉ lệ fail mà bỏ độ sâu queue (mục 6); cảnh báo không có người nhận nên "đặt" mà không bao giờ thấy.

  10. Load test: đẩy 10.000 job, đo throughput và queue depth theo thời gian; giết worker giữa chừng và chứng minh không mất job, không trùng tác dụng phụ (job có thể chạy lại, nên đo cả số lần chạy lại).

    Lời giải và cách kiểm tra

    Hướng làm: job ghi hai dấu vết khác nhau: bảng runs (mỗi lần chạy một dòng, để đếm chạy lại), effect_naive (luôn INSERT, mô phỏng tác dụng phụ không idempotent) và effect_idem (khoá chính, mô phỏng tác dụng phụ idempotent). SIGKILL worker khi đã xong khoảng 40%, bật worker mới, đợi đủ 10.000. Để thử nhanh, hạ lockDuration/stalledInterval xuống (mặc định 30 giây: job active của worker bị giết có thể phải đợi khoảng một phút mới được giao lại, suy luận từ mục 5, chưa đo ở cấu hình mặc định). Sơ đồ:

    textReady
     t0            ~40%: SIGKILL           +0,5 s              xong │ worker A ────────X                   │ worker B ──────────────►│ │ 20 job đang active (concurrency 20)  │                          │ │        │ khoá hết hạn sau lockDuration (5 s)                    │ │        └── stalled check (2 s) ► 20 job về `wait` ► B chạy lại ─┘ runs = 10.000 + 20 | effect_naive = 10.020 | effect_idem = 10.000

    Code tham chiếu (đã chạy):

    typescriptReady
    // worker (process con): src/kill-worker.tsconst worker = new Worker<{ n: number }>('load', async (job) => {  const n = job.data.n  await pool.query('INSERT INTO runs (n, pid) VALUES ($1, $2)', [n, process.pid])  await pool.query('INSERT INTO effect_naive (n) VALUES ($1)', [n])  await pool.query('INSERT INTO effect_idem (n) VALUES ($1) ON CONFLICT DO NOTHING', [n])  await new Promise((r) => setTimeout(r, 60))   // việc sau tác dụng phụ: cửa sổ để bị kill}, { connection, concurrency: 20, lockDuration: 5_000, stalledInterval: 2_000 })// driver: src/kill-test.ts (rút gọn)await queue.addBulk(Array.from({ length: TOTAL }, (_, n) => ({  name: 'work', data: { n }, opts: { jobId: `job-${n}`, attempts: 3, removeOnComplete: false },})))let w = start()                                 // spawn worker// đợi count(effect_idem) >= 40% TOTALw.kill('SIGKILL')await sleep(500)w = start()// đợi completed + failed >= TOTAL, rồi in://   runs, distinct n, reruns = runs - distinct, effect_naive, effect_idem, lost = TOTAL - effect_idem

    Kết quả mong đợi (đã chạy hai lần, Node 24.21, BullMQ 6.3.11, Redis 8.6.1, PostgreSQL 17.9 cục bộ; số giây chỉ mang tính tham khảo máy này):

    textReady
    {"TOTAL":10000,"secs":"34.2","completed":10000,"failed":0, "runs":10020,"distinctRuns":10000,"reruns":20, "effect_naive":10020,"naive_duplicates":20,"effect_idem":10000,"lost":0}

    Lần chạy thứ hai: secs 33.9, cùng các số còn lại. Đọc kết quả: không mất job (lost 0, completed 10000); số lần chạy lại bằng số job đang active lúc bị giết (20 = concurrency), tác dụng phụ ngây thơ bị trùng đúng 20, tác dụng phụ idempotent không trùng. Job chỉ chạy lại khi nó đang active lúc worker chết, nên reruns tối đa bằng concurrency; con số cụ thể phụ thuộc thời điểm giết và sẽ khác giữa các lần chạy. Lỗi hay gặp: giết bằng SIGTERM (worker tự close() nên không chạy lại, không chứng minh được gì); đợi theo mặc định 30 giây và tưởng job mất; đo "không trùng" bằng count(*) rồi quên rằng bảng ngây thơ cho ra số lớn hơn 10.000; không đặt jobId nên chạy lại bài test sẽ nhân đôi job; quên queue.obliterate giữa các lần đo.

Mục 10 là sản phẩm giao quan trọng nhất. "Tôi đã giết worker giữa lúc chạy 10.000 job và chứng minh được không mất job, không trùng tác dụng phụ" là một câu chuyện mạnh hơn mọi mô tả kiến trúc.

Khung và mã dùng chung

Tự làm trước, rồi mới mở. Môi trường đã chạy cho mọi mục có nhãn "Đã chạy": Node 24.21, BullMQ 6.3.11 (cài kèm ioredis), Redis 8.6.1 với maxmemory-policy noeviction, PostgreSQL 17.9, Vitest 5.0, TypeScript 7.0 (tsc --noEmit sạch). Mọi thứ chạy trên máy tạm, không dùng Docker và không dùng Prisma (SQL chạy bằng driver pg).

Sơ đồ tổng. Mọi yêu cầu ở trên là một ô trong hình này.

textReady
                    ┌────────── api (N instance) ──────────┐ client ─HTTP─►     │ 1 transaction: ghi dữ liệu + outbox  │                    └──────────────┬───────────────────────┘                                   │ PostgreSQL              ┌────────────────────▼─────────────────────┐              │ worker (process riêng, cùng image)       │              │  relay (SKIP LOCKED) ──► queue.add       │              └─────┬────────────────────┬───────────────┘                    ▼                    ▼   Redis(queue, noeviction)     idempotency_key (PENDING→DONE)   critical │ default │ bulk          │        │ job ─ retry+jitter ─► handler ─► provider (cùng key)        │                       │        │ hết retry / lỗi vĩnh viễn        ▼       DLQ ◄── /admin/dlq (xem, chạy lại có chủ đích)

Done khi#

  • Giải thích được 3 lợi ích của xử lý nền và cái giá phải trả

    Đáp án

    Độ trễ (trả 202 trong vài chục ms thay vì chờ việc nặng), chống chịu (worker chết thì job vẫn nằm trong queue), điều tiết tải (spike được queue giữ lại, worker rút theo tốc độ bạn chọn). Cái giá: hệ thống bất đồng bộ khó debug (job chạy hai lần, sai thứ tự, trễ, hoặc im lặng không chạy). Sai thường gặp: chỉ kể lợi ích. Xem GĐ10 mục 1.

  • Phân biệt queue vs log (BullMQ vs Kafka); biết khi nào cần replay

    Đáp án

    Queue (BullMQ, SQS): job bị tiêu thụ rồi biến mất, một job một worker. Log (Kafka): sự kiện được giữ lại, mỗi consumer group có offset riêng và replay được. Cần "phát lại 3 ngày sự kiện cho service mới" thì cần log; "gửi email này một lần" thì queue đủ. Sai thường gặp: chọn Kafka cho monolith vì "nó mạnh hơn". Xem GĐ10 mục 2.

  • Giải thích vì sao exactly-once không tồn tại và đạt hiệu quả đó bằng idempotency

    Đáp án

    Worker làm xong, ACK mất hoặc worker chết trước khi ACK tới: broker không phân biệt "chết trước khi làm" với "chết sau khi làm", nên chỉ chọn được giao lại (lặp) hoặc bỏ qua (mất). Thực tế là at-least-once, cộng handler idempotent để kết quả như exactly-once. Sai thường gặp: nói "dùng Kafka/SQS FIFO là có exactly-once" mà không nêu giới hạn. Xem GĐ10 mục 3 (có sơ đồ).

  • Viết được handler idempotent bằng cả 3 cách ở mục 3

    Đáp án

    A: PENDING → DONE có locked_until và lock_token trong DB (không ghi khoá rồi mới làm: lỗi tạm sau khi giữ khoá sẽ mất giao dịch). B: thao tác tự nhiên idempotent (gán, ghi đè cùng key). C: đẩy khoá sang bên thứ ba (idempotencyKey của Stripe). Cách tự kiểm: test gọi hai lần cho đúng một tác dụng phụ, và test lỗi tạm sau khi giành khoá (bài tập mục 11, yêu cầu 4). Xem GĐ10 mục 3.

  • Giải thích outbox và vì sao queue.add() trong transaction sai theo cả hai chiều

    Đáp án

    queue.add() trong transaction sai theo chiều 1 (add xong rồi rollback: job cho đơn không tồn tại) và chiều 2 (commit xong rồi Redis chết: đơn có mà không ai được báo). Outbox ghi ý định vào cùng DB cùng transaction, relay đẩy sang queue sau. Vẫn là at-least-once nên handler vẫn phải idempotent. Xem GĐ10 mục 4 (có sơ đồ).

  • Hiểu FOR UPDATE SKIP LOCKED giải quyết gì và vì sao nó cho phép nhiều relay

    Đáp án

    Khoá các hàng được chọn, và relay thứ hai bỏ qua hàng đang bị khoá thay vì chờ, nên nhiều relay chạy song song mà không giẫm nhau. Không có nó, chạy nhiều relay chỉ xếp hàng chờ nhau. Sai thường gặp: nghĩ nó chống trùng sau khi relay chết (không: hàng được nhả khi rollback, nên vẫn có thể add lại). Xem GĐ10 mục 4.

  • Payload nhỏ, có v; biết vì sao không nhét cả object và không để PII trong Redis

    Đáp án

    Chỉ chứa id và phiên bản: payload là ảnh chụp lúc enqueue (có thể cũ khi worker chạy), nằm trong Redis nên PII là rủi ro tuân thủ, payload lớn tốn RAM Redis. v cho phép handler xử lý cả hình dạng cũ lẫn mới khi job cũ còn trong queue. Xem GĐ10 mục 5.

  • Retry có backoff + jitter; giải thích thundering herd

    Đáp án

    Backoff mũ giãn các lần thử; jitter ngẫu nhiên hoá để 5.000 job cùng fail không cùng retry đúng lúc, tạo thundering herd vào API vừa hồi phục. Ở BullMQ 6, jitter đặt trong backoff, khoảng chờ rơi trong [gốc × (1 − jitter), gốc] (đã chạy: gốc 300 ms cho 277 ms). Xem GĐ10 mục 5.

  • Phân biệt lỗi tạm thời / vĩnh viễn; không retry lỗi vĩnh viễn

    Đáp án

    Tạm: 429, 408, 5xx, timeout mạng, lỗi lock → để retry. Vĩnh viễn: 4xx cú pháp, email sai, record đã xoá → UnrecoverableError để vào failed ngay. Cách tự kiểm: job lỗi vĩnh viễn chỉ chạy đúng 1 lần (đã chạy) và vẫn vào DLQ (xem mục 6). Sai thường gặp: return êm nên job hiện completed.

  • Biết lockDuration do worker tự gia hạn, job chỉ stalled khi event loop bị chặn hoặc worker chết, và grace period phải lớn hơn job dài nhất

    Đáp án

    Worker tự gia hạn khoá mỗi lockRenewTime (mặc định nửa lockDuration), nên job dài hơn lockDuration vẫn an toàn. Job chỉ stalled khi worker không gia hạn được: event loop bị chặn quá lockDuration, hoặc worker chết/mất Redis; khi đó job chạy lại ở worker khác. Grace period phải lớn hơn job dài nhất (đã trừ preStop), còn lockDuration không nằm trong bất đẳng thức đó. Xem GĐ10 mục 5 và mục 7.

  • Worker chạy process riêng, xử lý SIGTERM, chờ job đang chạy xong

    Đáp án

    Process riêng để job nặng không chặn API, scale và deploy độc lập. Bắt SIGTERM: cờ shuttingDown, hẹn giờ thoát cưỡng bức nhỏ hơn grace period, await worker.close(), rồi mới đóng DB/Redis. Cách tự kiểm: kill -TERM lúc job đang chạy, phải thấy job hoàn tất và exit=0 (đã chạy: close() mất 821 ms). Xem GĐ10 mục 7 (có sơ đồ).

  • Tính được concurrency tối đa từ connection pool

    Đáp án

    Mỗi job giữ một connection khi đang truy vấn nên cần concurrency của một process không vượt pool của process đó, và số process × pool không vượt giới hạn kết nối của DB (còn phần của API). Ví dụ: 20 worker × concurrency 10 = 200 kết nối cùng lúc; với Postgres cho phép vài trăm kết nối thì phải hạ concurrency hoặc đặt PgBouncer. Việc CPU-bound thì concurrency lớn hơn số core là vô ích. Xem GĐ10 mục 7, Pitfall #5.

  • Có DLQ; job không bao giờ biến mất im lặng; chạy lại là hành động có chủ đích

    Đáp án

    Job hết retry hoặc lỗi vĩnh viễn đi vào queue riêng kèm lỗi và payload, có log và metric; chạy lại là việc của người sau khi đã sửa nguyên nhân và kiểm tra handler idempotent. Cách tự kiểm: làm một job fail cố ý, thấy dlqSize tăng 1 và có dòng log job dead-lettered. Sai thường gặp: điều kiện attemptsMade >= attempts bỏ sót lỗi vĩnh viễn (đã chạy, xem mục 6). Xem GĐ10 mục 6.

  • Theo dõi queue depth và oldest job age, không chỉ failure rate

    Đáp án

    Độ sâu queue tăng đơn điệu nghĩa là throughput worker thấp hơn enqueue, tín hiệu sớm trước khi có job nào fail; tuổi job cũ nhất đối chiếu với SLA (email trễ quá 5 phút). Cách tự kiểm: dừng worker, enqueue 5 job, thấy depth 5 và oldestAgeMs tăng dần (đã chạy). Xem GĐ10 mục 6.

  • Chạy cron đúng một lần trên N instance; biết giới hạn của khoá Redis và vì sao vẫn cần idempotency

    Đáp án

    Dùng BullMQ scheduler (Redis điều phối), K8s CronJob hoặc advisory lock Postgres. Khoá Redis (SET NX PX) chỉ giảm việc trùng: process treo quá TTL thì instance khác lấy được khoá và hai bên cùng chạy, nên tính đúng đắn vẫn dựa vào idempotency. Sai thường gặp: setInterval trong app rồi scale ra 3 instance. Xem GĐ10 mục 8 và 8.1.

  • Cron có tz tường minh; xử lý được missed run

    Đáp án

    tz luôn tường minh (ví dụ Asia/Ho_Chi_Minh); cron ở múi giờ có DST có thể bỏ hoặc chạy hai lần. Missed run chọn chính sách có chủ đích: bỏ qua, chạy bù một lần, hầu như không bao giờ chạy bù tất cả. Cách tự kiểm: getJobSchedulers() lúc boot không tăng sau deploy. Xem GĐ10 mục 8.2.

  • Thiết kế được job theo con trỏ dữ liệu thay vì theo thời điểm

    Đáp án

    Lưu lastProcessedAt, mỗi lần chạy xử lý các ngày đã kết thúc sau con trỏ, ghi đè theo ngày, rồi mới dịch con trỏ. Tự bù missed run và chạy lại vô hại. Cách tự kiểm: test con trỏ 1/10, now 5/10 xử lý ngày 2, 3, 4; chạy lại xử lý 0 (đã chạy, bài tập mục 11, yêu cầu 8). Xem GĐ10 mục 8.2.

  • Chống được noisy neighbor bằng queue riêng hoặc rate limit theo tenant

    Đáp án

    Queue riêng theo mức ưu tiên (critical so với bulk) làm trước; priority có nguy cơ starvation; rate limit theo tenant (BullMQ Pro có group, bản OSS dùng token bucket đầu handler rồi re-enqueue có delay). Cách tự kiểm: tenant A đẩy hàng loạt, job của tenant B vẫn bắt đầu trong thời gian ngắn. Xem GĐ10 mục 9.

  • Test handler không cần Redis; test tích hợp không dùng sleep cố định

    Đáp án

    Handler là hàm nhận payload và deps, test gọi hai lần rồi expect(send).toHaveBeenCalledTimes(1); test tích hợp chờ bằng QueueEvents hoặc polling có timeout, không sleep cố định. Cách tự kiểm: test handler chạy được khi tắt Redis. Xem GĐ10 mục 10.

  • Chọn đúng cấp khoá advisory của Postgres khi đi qua connection pool

    Đáp án

    Khoá cấp transaction: pg_try_advisory_xact_lock trong BEGIN ... COMMIT, tự nhả khi kết thúc transaction hoặc mất connection. Khoá cấp session gắn với connection mà connection trong pool sống lâu: khoá có thể còn lại sau khi trả connection, và pg_advisory_unlock ở connection khác trả false kèm cảnh báo. Cách tự kiểm: ba instance gọi cùng lúc, đúng một instance chạy, hai instance bỏ qua. Sai thường gặp: giữ transaction cho việc dài (hết timeout của Prisma, giữ connection lâu). Xem GĐ10 mục 8.1.

  • Chờ job trong test bằng waitUntilFinished, không bằng sleep hay listener tạo vội

    Đáp án

    Một QueueEvents dùng chung, await events.waitUntilReady() trước khi add, await job.waitUntilFinished(events, ttl) (reject khi job fail hoặc quá ttl), events.close() ở cuối. Cách tự kiểm: một test job thành công và một test job UnrecoverableError phải đỏ/xanh đúng như mong đợi. Sai thường gặp: new QueueEvents(...).on('completed', ...) tạo ngay trước khi add (bỏ lỡ sự kiện, resolve theo job bất kỳ). Xem GĐ10 mục 10.

  • Dùng được delay, flow cha-con và limiter, và nói được mỗi cái không phải là gì

    Đáp án

    delay hoãn một job nhưng không bảo đảm đúng giờ và không phải lịch (lịch là scheduler). Flow: cha chạy khi mọi con xong, getChildrenValues() gom kết quả; con lỗi thì cha kẹt ở waiting-children trừ khi đặt failParentOnFailure hoặc ignoreDependencyOnFailure; flow không phải saga (không bồi hoàn). limiter giới hạn cả queue, không theo tenant. Cách tự kiểm: hai worker với limiter: { max: 3, duration: 1000 } xử lý 9 job theo 3, 3, 3 mỗi giây. Xem GĐ10 mục 5.1.

  • Có checklist vận hành BullMQ: removeOn*, kết nối producer/worker, on('error')

    Đáp án

    defaultJobOptions có removeOnComplete/removeOnFail (mặc định job nằm lại mãi); producer enableOfflineQueue: false để lỗi nhanh khi Redis chết giữa chừng; worker maxRetriesPerRequest: null; on('error') cho Worker và Queue; bắt unhandledRejection/uncaughtException. Cách tự kiểm: tắt Redis rồi gọi queue.add từ API, phải nhận lỗi ngay chứ không treo. Xem GĐ10 mục 7.1.


Câu hỏi mở / chưa giải quyết#

  • Khi nào chuyển từ BullMQ sang Kafka? Tín hiệu: cần replay lịch sử, cần nhiều consumer group độc lập trên cùng luồng sự kiện, hoặc throughput vượt khả năng của một Redis. Không phải "vì hệ thống lớn".

    Hướng trả lời hiện tại

    Hướng trả lời hiện tại (chưa khẳng định): hỏi lại ba điều trước khi chuyển. Có cần đọc lại sự kiện cũ không? Có hơn một nhóm tiêu thụ độc lập trên cùng luồng không? Số liệu đo (không phải ước đoán) cho thấy một Redis không đủ không? Nếu cả ba đều "không", ở lại BullMQ. Chưa có ngưỡng số cụ thể trong tài liệu này; hãy đo trên tải của bạn.

  • pg-boss có đủ thay BullMQ không? Với tải vừa và khi bạn đã có Postgres, có — và bạn được thêm lợi ích lớn là enqueue trong cùng transaction với dữ liệu, tức là không cần outbox. Đánh đổi: throughput thấp hơn Redis và tạo tải ghi lên DB chính.

    Hướng trả lời hiện tại

    Hướng trả lời hiện tại (chưa khẳng định): với một monolith đã có Postgres và tải vừa, pg-boss đáng thử trước vì bỏ được Redis và outbox. Ý "enqueue trong cùng transaction" dựa trên mô tả của thư viện, chưa chạy ở đây; kiểm tra trên phiên bản bạn dùng trước khi dựa vào nó. Cần dữ liệu thật về throughput trước khi kết luận thay hẳn.

  • Saga cho quy trình dài nhiều bước (đặt hàng → thanh toán → kho → giao hàng, có bồi hoàn khi hỏng giữa chừng) thuộc GĐ20, không thuộc giai đoạn này.

    Hướng trả lời hiện tại

    Hướng trả lời hiện tại (chưa khẳng định): trong phạm vi GĐ10, mỗi bước của saga là một job idempotent có bước bồi hoàn tương ứng, nối với nhau qua outbox; phần điều phối và mô hình hoá thuộc GĐ20.