Chuyển tới nội dung chính

17.14 — Ví dụ thực tế nhanh

Tóm tắt

Một tình huống ngắn nhưng rất hay gặp: queue lead-converted bình thường luôn rỗng, rồi đột ngột dồn 400.000 message trong 20 phút. Không có lỗi nào trong log, consumer vẫn "chạy", và CPU của nó chỉ 8%. Bài này cho thấy cách chẩn đoán loại sự cố này bằng đúng ba chỉ số — độ sâu queue, số message chưa ACK, và số consumer đang gắn — và bảng đối chiếu ba chỉ số đó với năm nguyên nhân phổ biến. Điều đáng nhớ nhất: queue dồn gần như không bao giờ là lỗi của broker — nó là triệu chứng của consumer, và ba chỉ số này chỉ thẳng ra consumer đang sai ở đâu.

Mục tiêu bài học​

Sau bài này bạn có thể:

  • Chẩn đoán queue dồn bằng ba chỉ số.
  • Phân biệt năm nguyên nhân phổ biến.
  • Xử lý sự cố ngay và phòng ngừa sau đó.
  • Đặt cảnh báo đúng cho hệ thống message.

Nội dung bài học​

17.14.1 — Tình huống​

14:00  Queue "lead-converted": 0 message (binh thuong)
14:05 8.000 message
14:12 120.000 message
14:20 400.000 message va van tang
Consumer: đang chạy, CPU 8%, KHÔNG có lỗi trong log

Bộ nhớ RabbitMQ tăng dần, và nếu chạm ngưỡng cao thì broker sẽ chặn publisher — nghĩa là API nghiệp vụ cũng dừng theo.

17.14.2 — Ba chỉ số cần xem​

rabbitmqctl list_queues name messages messages_unacknowledged consumers
name             messages  messages_unacknowledged  consumers
lead-converted 400123 1 1

Đọc ba con số cuối:

Chỉ sốÝ nghĩaGiá trị ở đây
messagesTổng đang chờ400.123 — dồn nặng
messages_unacknowledgedĐang được xử lý1 — chỉ xử lý một message một lúc
consumersSố consumer đang gắn1 — chỉ một instance

Hai con số cuối cho câu trả lời ngay: một consumer, xử lý tuần tự từng message một. Với 400.000 message và mỗi message tốn 50ms, cần 5,5 giờ để rút hết.

17.14.3 — Bảng chẩn đoán​

messagesunackedconsumersNguyên nhân
Cao00Consumer đã chết hoặc mất kết nối
CaoThấpThấpThiếu consumer / prefetch quá nhỏ
CaoCaoBình thườngConsumer chậm — xử lý lâu hoặc chờ I/O
CaoDao động, message quay lạiBình thườngPoison message — lỗi rồi requeue vô hạn
Thấp nhưng tăng dầnBình thườngBình thườngTốc độ publish vượt tốc độ xử lý

Ca này khớp hàng thứ hai. Hàng thứ tư đáng cảnh giác nhất: nhìn qua trông như consumer đang làm việc, nhưng thực ra nó xử lý cùng một message mãi mãi (bài 17.6).

17.14.4 — Nguyên nhân gốc​

// Cấu hình hiện tại
cfg.ReceiveEndpoint("lead-converted", e =>
{
e.PrefetchCount = 1; // chỉ giữ 1 message một lúc
e.ConcurrentMessageLimit = 1; // chỉ xử lý 1 message một lúc
e.Consumer<LeadConvertedConsumer>();
});
replicas: 1                           # chi MOT instance

Cấu hình này giới hạn thông lượng ở một message tại một thời điểm, trên một instance. Nó đủ dùng khi lưu lượng bình thường (vài chục message mỗi phút), nhưng một chiến dịch marketing đẩy 400.000 lead vào hệ thống thì không cách nào theo kịp.

Vì sao PrefetchCount = 1? Vì ai đó từng gặp vấn đề phân bổ tải không đều và đặt nó xuống 1 để "cho chắc" — một cách sửa quá tay (bài 17.4).

17.14.5 — Xử lý ngay​

cfg.ReceiveEndpoint("lead-converted", e =>
{
e.PrefetchCount = 32; // giữ sẵn 32 message trong bộ đệm
e.ConcurrentMessageLimit = 8; // xử lý 8 message song song
e.Consumer<LeadConvertedConsumer>();
});
kubectl scale deployment lead-consumer --replicas=6
Thông lượng: 20 msg/s -> 960 msg/s
Thời gian rút hết 400.000 message: 5,5 giờ -> ~7 phút

Điều kiện tiên quyết trước khi tăng song song: consumer phải idempotent. Xử lý song song làm tăng khả năng hai message liên quan chạy đồng thời và va vào nhau (bài 17.2).

Và phải kiểm tra database chịu được không: 6 instance × 8 luồng = 48 truy vấn song song. Nếu connection pool hoặc database không chịu nổi, bạn chỉ chuyển nghẽn từ queue sang database (bài 19.7).

17.14.6 — Phòng ngừa​

1. Cảnh báo theo xu hướng, không theo giá trị tuyệt đối:

CẢNH BÁO khi: độ sâu queue tăng liên tục trong 5 phút
CẢNH BÁO khi: tuổi message cũ nhất > 2 phút
CẢNH BÁO khi: số consumer = 0

Ngưỡng tuyệt đối như "queue > 10.000" gây báo động giả vào lúc cao điểm bình thường và im lặng khi một consumer duy nhất chết trong lúc lưu lượng thấp.

Chỉ số tuổi message cũ nhất tốt hơn số lượng: 400.000 message vừa được publish trong một giây là bình thường; một message kẹt 10 phút là sự cố (bài 17.8).

2. Tự co giãn theo độ sâu queue, không theo CPU:

# KEDA — scale theo so message dang cho
triggers:
- type: rabbitmq
metadata:
queueName: lead-converted
queueLength: "50" # moi 50 message them mot pod

CPU là chỉ số sai ở đây: consumer chờ I/O có CPU rất thấp ngay cả khi queue dồn 400.000 message (bài 19.4).

3. Test tải cho consumer. Đẩy 100.000 message vào staging và đo thông lượng thật. Con số đó cho biết hệ thống chịu được đỉnh tải cỡ nào trước khi đỉnh tải xảy ra.

17.14.7 — Rà lại code của bạn​

Danh sách rà soát queue dồn

  • •PrefetchCount không bị đặt bằng 1 mà không có lý do rõ ràng.
  • •ConcurrentMessageLimit đặt theo khả năng thật của consumer.
  • •Consumer idempotent trước khi tăng mức song song.
  • •Đã kiểm tra database chịu được số truy vấn song song sau khi scale.
  • •Có cảnh báo khi số consumer bằng 0.
  • •Cảnh báo dựa trên xu hướng và tuổi message, không chỉ số lượng.
  • •Tự co giãn theo độ sâu queue, không theo CPU.
  • •Đã đo thông lượng thật của consumer bằng test tải.
  • •Có DLQ và cảnh báo trên DLQ để bắt poison message.

Bài tập áp dụng​

Bài 1 — Đọc ba chỉ số​

Chạy rabbitmqctl list_queues trên môi trường của bạn và đối chiếu với bảng chẩn đoán ở mục 17.14.3.

Tiêu chí hoàn thành: bạn đọc được ba chỉ số cho mọi queue, và phân biệt được tình trạng bình thường với tình trạng cần hành động cho từng loại consumer.

Gợi ý và lời giải — Bài 1

Gợi ý. unacked = 32 là bình thường hay bất thường? Câu trả lời phụ thuộc vào PrefetchCount.

Lời giải:

rabbitmqctl list_queues name messages messages_ready messages_unacknowledged \
consumers consumer_utilisation memory
name                     messages  ready  unacked  consumers  utilisation  memory
lead-converted 412048 412016 32 1 0.98 84213504
order-created 18 2 16 4 0.45 1204224
email-notifications 0 0 0 0 0.00 204800
erp-sync 847 843 4 2 1.00 2048576
audit-log 12904 12904 0 3 0.00 4194304

Đọc từng dòng:

lead-converted — 412.048 message, 1 consumer, utilisation 0,98:

messages cao + consumers = 1 + utilisation ≈ 1
-> consumer đang chạy HẾT CÔNG SUẤT nhưng vẫn không kịp
-> CẦN: thêm consumer, tăng đồng thời, hoặc tối ưu xử lý

consumer_utilisation gần 1 nghĩa là consumer gần như không bao giờ rảnh — nó là chỉ số tốt nhất để phân biệt "consumer chậm" với "thiếu consumer".

order-created — 18 message, 4 consumer, utilisation 0,45:

messages thấp + utilisation 0,45
-> consumer RẢNH gần một nửa thời gian
-> bình thường, thậm chí có thể GIẢM số consumer

email-notifications — mọi thứ bằng 0:

consumers = 0 -> KHÔNG có consumer nào kết nối

Đây là tình trạng nguy hiểm nhất vì nó trông bình thường: messages = 0 nên dashboard xanh, không có cảnh báo nào. Nhưng nếu có message tới, không ai xử lý.

rabbitmqctl list_consumers | grep email-notifications
(không có kết quả)

Nguyên nhân thường gặp: dịch vụ không chạy, sai chuỗi kết nối, hoặc consumer crash lúc khởi động và pod đang ở CrashLoopBackOff.

erp-sync — utilisation = 1,00 chính xác:

utilisation = 1,00 và unacked = 4 với 2 consumer
-> mỗi consumer giữ đúng 2 message, và không bao giờ rảnh
-> PrefetchCount có thể quá thấp

audit-log — 12.904 message, 3 consumer, unacked = 0, utilisation = 0:

Có consumer, nhưng KHÔNG nhận message nào
-> consumer đăng ký nhưng không lấy message
-> thường do: channel bị chặn (flow control), hoặc consumer treo
rabbitmqctl list_channels name state messages_unacknowledged
name                      state    unacked
<rabbit@node.1.2345.0> flow 0 <- FLOW CONTROL

Trạng thái flow nghĩa là broker đang chặn channel vì quá tải — thường do bộ nhớ hoặc đĩa chạm ngưỡng:

rabbitmqctl status | grep -A5 "memory\|disk_free"

Phân biệt bình thường và cần hành động — phụ thuộc vào cấu hình:

unacked ≈ PrefetchCount × consumers   -> BÌNH THƯỜNG, consumer đang bận đúng mức
unacked = 0 với messages > 0 -> BẤT THƯỜNG, consumer không lấy message
unacked = PrefetchCount × consumers và đứng yên nhiều phút -> consumer TREO
lead-converted:  prefetch 32 × 1 consumer = 32 -> unacked 32 là ĐÚNG mức
order-created: prefetch 4 × 4 consumer = 16 -> unacked 16 là ĐÚNG mức
audit-log: prefetch 16 × 3 consumer = 48 -> unacked 0 là SAI

Nghĩa là bạn không thể đọc unacked mà không biết PrefetchCount. Ghi nó vào tài liệu vận hành cho từng queue.

Bảng chẩn đoán mở rộng:

messagesunackedconsumersutilisationChẩn đoán
Cao00—Consumer chết hoặc chưa khởi động
Cao0> 00Channel bị flow control hoặc consumer treo
Cao≈ prefetch × n> 0≈ 1Thiếu năng lực — thêm consumer
CaoThấp> 0ThấpPrefetch quá nhỏ
CaoDao động, messages không giảm> 0≈ 1Poison message quay vòng
Thấp, tăng dầnBình thường> 00,7–0,9Sắp không kịp — theo dõi
Thấp, ổn địnhBình thường> 0< 0,5Bình thường

Ba lệnh bổ sung khi chẩn đoán:

# 1. Tốc độ vào và ra
rabbitmqctl list_queues name messages message_stats.publish_details.rate \
message_stats.deliver_get_details.rate
name              messages  publish_rate  deliver_rate
lead-converted 412048 1840.2 112.4
publish 1.840/s, deliver 112/s
-> vào nhanh hơn ra 16 lần
-> hàng đợi sẽ tiếp tục tăng
-> cần năng lực gấp 16 lần, hoặc giảm tốc độ publish

Hai con số này cho bạn ước lượng thời gian — thứ cần nhất khi quyết định "chờ hay can thiệp":

412.048 / (112,4 - 1.840,2) -> số âm -> KHÔNG BAO GIỜ rút hết
# 2. Message ở dead-letter queue
rabbitmqctl list_queues name messages | grep -i "dlq\|dead\|error\|_skipped"
lead-converted_error    1847

1.847 message thất bại vĩnh viễn — một con số mà không ai biết cho tới khi có người nhìn.

# 3. Message cũ nhất trong queue
rabbitmqadmin get queue=lead-converted count=1 ackmode=reject_requeue_true
timestamp: 1758787200        <- xem message đầu hàng đợi tới từ bao giờ

Và chỉ số quan trọng nhất không có trong list_queues: tuổi của message cũ nhất.

public async Task Consume(ConsumeContext<LeadConvertedV1> ctx)
{
var tuoi = DateTime.UtcNow - (ctx.SentTime ?? DateTime.UtcNow);
_doTre.Record(tuoi.TotalSeconds,
new KeyValuePair<string, object?>("queue", "lead-converted"));

// ...
}
Cảnh báo khi p95 của độ trễ vượt SLA của luồng đó.

Như đã nói ở bài 17.7: số lượng nói về tải, tuổi nói về sự cố. 412.048 message vừa được đẩy vào trong một phút là bình thường; 412.048 message trong đó cái cũ nhất từ ba ngày trước là một sự cố.


Bài 2 — Đo thông lượng​

Đẩy 50.000 message vào staging, đo thời gian rút hết, rồi tăng PrefetchCount và ConcurrentMessageLimit và đo lại.

Tiêu chí hoàn thành: bạn tìm được điểm bão hoà, và biết nút thắt thật nằm ở đâu.

Gợi ý và lời giải — Bài 2

Gợi ý. Tăng đồng thời mãi có tăng thông lượng mãi không? Cái gì chặn lại trước?

Lời giải — đo có hệ thống:

[Theory]
[InlineData(1, 1)]
[InlineData(8, 4)]
[InlineData(16, 8)]
[InlineData(32, 16)]
[InlineData(64, 32)]
[InlineData(128, 64)]
public async Task Do_thong_luong(int prefetch, int dongThoi)
{
await XoaQueueAsync();
await DayMessageAsync(50_000);

var cfg = TaoCauHinh(prefetch, dongThoi);
var sw = Stopwatch.StartNew();

await ChayConsumerToiKhiRutHetAsync(cfg);

sw.Stop();
_output.WriteLine($"prefetch={prefetch,3} đồng thời={dongThoi,3} " +
$"{sw.Elapsed.TotalSeconds,6:F1}s " +
$"{50_000 / sw.Elapsed.TotalSeconds,7:F0} msg/s");
}
prefetch=  1 đồng thời=  1   412.3s     121 msg/s
prefetch= 8 đồng thời= 4 104.8s 477 msg/s
prefetch= 16 đồng thời= 8 54.2s 923 msg/s
prefetch= 32 đồng thời= 16 31.7s 1,577 msg/s
prefetch= 64 đồng thời= 32 28.4s 1,761 msg/s
prefetch=128 đồng thời= 64 29.1s 1,718 msg/s <- CHẬM HƠN

Điểm bão hoà nằm ở khoảng đồng thời 32. Vượt qua đó, thông lượng giảm.

1 -> 4:    nhanh gấp 3,9 lần   (gần tuyến tính)
4 -> 8: nhanh gấp 1,9 lần
8 -> 16: nhanh gấp 1,7 lần
16 -> 32: nhanh gấp 1,1 lần <- lợi ích giảm mạnh
32 -> 64: CHẬM đi 2% <- vượt điểm bão hoà

Tìm nút thắt thật:

# CPU của consumer
kubectl top pods -l app=lead-consumer
NAME                CPU(cores)   MEMORY(bytes)
lead-consumer-abc 180m 412Mi
CPU 180m trên giới hạn 1000m -> chỉ 18%
-> KHÔNG phải nút thắt CPU
-- Database
SELECT TOP 10
wait_type, wait_time_ms, waiting_tasks_count,
wait_time_ms / NULLIF(waiting_tasks_count, 0) AS avg_wait_ms
FROM sys.dm_os_wait_stats
WHERE wait_type NOT IN ('CLR_SEMAPHORE','SLEEP_TASK','BROKER_TASK_STOP','XE_TIMER_EVENT')
ORDER BY wait_time_ms DESC;
wait_type                wait_time_ms  waiting_tasks  avg_wait_ms
PAGEIOLATCH_SH 1842104 84210 21
WRITELOG 412048 98204 4
LCK_M_U 184210 12048 15
PAGEIOLATCH_SH cao -> chờ đọc trang từ đĩa
-> DATABASE là nút thắt, không phải consumer
-- Số kết nối đang dùng
SELECT COUNT(*) AS SoKetNoi, DB_NAME(database_id) AS Db
FROM sys.dm_exec_sessions WHERE is_user_process = 1
GROUP BY database_id;
SoKetNoi  Db
100 Crm <- CHẠM TRẦN connection pool (mặc định 100)

Đây là nút thắt thật. Tăng đồng thời lên 64 nghĩa là 64 consumer cùng cần kết nối, nhưng pool chỉ có 100 và chúng phải chia sẻ với API — nên chúng chờ nhau.

Timeout expired. The timeout period elapsed prior to obtaining a connection
from the pool.

Ba nút thắt thường gặp, theo thứ tự phổ biến:

Nút thắtDấu hiệuCách xác nhận
DatabasePAGEIOLATCH, WRITELOG, LCK_M_* caodm_os_wait_stats
Connection poolTimeout lấy kết nốiĐếm session, xem Max Pool Size
API bên ngoàiConsumer chờ I/O, CPU thấpTrace, đo thời gian từng chặng
Brokerconsumer_utilisation thấp dù nhiều messagerabbitmqctl status
CPU consumerCPU chạm giới hạnkubectl top

Dòng đầu và dòng ba chiếm phần lớn trường hợp, và chúng có một điểm chung: tăng đồng thời làm chúng tệ hơn, không tốt hơn.

Database là nút thắt -> thêm consumer = thêm truy vấn đồng thời
-> nhiều khoá hơn, nhiều deadlock hơn
-> thông lượng GIẢM

Sửa đúng nút thắt:

// 1. Tăng connection pool — nhưng cẩn thận
options.UseSqlServer(conn + ";Max Pool Size=200;");
CẢNH BÁO: mỗi kết nối tốn bộ nhớ ở phía database.
200 kết nối × 3 pod = 600 kết nối -> có thể vượt giới hạn của database.
Đo trước khi tăng.
// 2. Gom lô — cách hiệu quả nhất cho nút thắt database
cfg.ReceiveEndpoint("lead-converted", e =>
{
e.PrefetchCount = 100;
e.Batch<LeadConvertedV1>(b =>
{
b.MessageLimit = 50;
b.TimeLimit = TimeSpan.FromSeconds(2);
b.Consumer<LeadConvertedBatchConsumer>(context);
});
});
public class LeadConvertedBatchConsumer : IConsumer<Batch<LeadConvertedV1>>
{
public async Task Consume(ConsumeContext<Batch<LeadConvertedV1>> ctx)
{
var ids = ctx.Message.Select(m => m.Message.LeadId).ToList();

// MỘT truy vấn cho 50 message thay vì 50 truy vấn
var leads = await _db.Leads.Where(l => ids.Contains(l.Id)).ToListAsync(ct);

foreach (var msg in ctx.Message) { /* xử lý */ }

await _db.SaveChangesAsync(ct); // MỘT transaction
}
}
prefetch=100 batch=50   9.2s   5,435 msg/s      <- nhanh gấp 3 lần điểm bão hoà cũ

Gom lô giảm số lần đi về database từ 50.000 xuống 1.000 — và đó là lý do nó vượt qua giới hạn mà tăng đồng thời không vượt được.

Và một điều cần biết về batch: xử lý lỗi phức tạp hơn.

Một message trong lô thất bại -> cả lô bị retry?
-> phải tách message lỗi ra, hoặc chấp nhận xử lý lại 49 message kia
-> và consumer PHẢI idempotent (bài 17.2) — điều này giờ là bắt buộc

Quy trình tìm cấu hình đúng:

1. Đo với cấu hình hiện tại -> có con số cơ sở
2. Tăng dần đồng thời tới khi thông lượng NGỪNG tăng
3. Tìm nút thắt tại điểm đó (database? pool? API ngoài?)
4. Sửa nút thắt
5. Quay lại bước 2

Dừng khi: thông lượng đủ cho nhu cầu, hoặc nút thắt tiếp theo quá đắt để sửa.

Và một lời nhắc: đừng tối ưu thông lượng khi vấn đề là độ trễ.

50.000 message trong 9 giây -> thông lượng tốt
Nhưng nếu message đầu tiên chờ 9 giây trước khi được xử lý,
và nghiệp vụ cần nó trong 1 giây -> thông lượng không giải quyết được gì

Hai vấn đề khác nhau, hai cách sửa khác nhau: thông lượng cần đồng thời và gom lô; độ trễ cần queue riêng với ưu tiên cao cho luồng cần nhanh.


Bài 3 — Tái hiện poison message​

Gửi một message gây exception và quan sát nó quay lại queue liên tục nếu không có DLQ.

Tiêu chí hoàn thành: bạn quan sát được vòng lặp, và cấu hình được DLQ cùng quy trình xử lý message trong đó.

Gợi ý và lời giải — Bài 3

Gợi ý. Message gây exception -> nack -> requeue -> lại gây exception. Vòng lặp này dừng khi nào?

Lời giải — tái hiện:

public class LeadConvertedConsumer : IConsumer<LeadConvertedV1>
{
public async Task Consume(ConsumeContext<LeadConvertedV1> ctx)
{
var lead = await _db.Leads.FirstAsync(l => l.Id == ctx.Message.LeadId, ct);
// Message có LeadId không tồn tại -> InvalidOperationException
}
}
cfg.ReceiveEndpoint("lead-converted", e =>
{
// KHÔNG cấu hình retry, KHÔNG cấu hình DLQ
e.ConfigureConsumer<LeadConvertedConsumer>(context);
});
# Gửi một message với LeadId không tồn tại
[08:14:22] Xử lý LeadConvertedV1 abc-999
[08:14:22] InvalidOperationException: Sequence contains no elements
[08:14:22] Xử lý LeadConvertedV1 abc-999
[08:14:22] InvalidOperationException: Sequence contains no elements
[08:14:22] Xử lý LeadConvertedV1 abc-999
... (liên tục, hàng nghìn lần mỗi giây)
rabbitmqctl list_queues name messages messages_unacknowledged consumers
lead-converted    1    1    1

Một message, nhưng consumer chạy 100% CPU — nó xử lý cùng message đó liên tục, không bao giờ dừng.

Ba hậu quả:

1. CPU của consumer chạm trần -> message HỢP LỆ không được xử lý
2. Log đầy -> đĩa đầy (bài 15.3)
3. Database nhận hàng nghìn truy vấn mỗi giây cho cùng một id

Và điều tệ nhất: messages = 1 nên dashboard trông bình thường. Không có cảnh báo nào về độ sâu hàng đợi.

Cấu hình đúng — ba tầng:

cfg.ReceiveEndpoint("lead-converted", e =>
{
// Tầng 1 — retry NGAY cho lỗi thoáng qua
e.UseMessageRetry(r =>
{
r.Exponential(3,
minInterval: TimeSpan.FromSeconds(1),
maxInterval: TimeSpan.FromSeconds(30),
intervalDelta: TimeSpan.FromSeconds(5));

// Lỗi vĩnh viễn -> KHÔNG retry, đi thẳng vào DLQ
r.Ignore<ValidationException>();
r.Ignore<JsonException>();
r.Ignore<BusinessRuleException>();
});

// Tầng 2 — redelivery có độ trễ dài, cho lỗi cần thời gian hồi phục
e.UseScheduledRedelivery(r => r.Intervals(
TimeSpan.FromMinutes(5),
TimeSpan.FromMinutes(15),
TimeSpan.FromMinutes(60)));

// Tầng 3 — hết cách -> DLQ (MassTransit tự tạo queue _error)
e.ConfigureConsumer<LeadConvertedConsumer>(context);
});
Message lỗi:
-> retry ngay 3 lần (1s, 6s, 30s)
-> nếu vẫn lỗi: redelivery sau 5 phút, 15 phút, 60 phút
-> nếu vẫn lỗi: vào lead-converted_error
rabbitmqctl list_queues name messages | grep error
lead-converted_error    1

Vì sao cần cả ba tầng:

Tầng 1 (retry ngay):     lỗi kéo dài vài giây — deadlock, timeout mạng
Tầng 2 (redelivery dài): lỗi kéo dài vài phút — dịch vụ đang restart,
database đang failover, dữ liệu chưa kịp đồng bộ
Tầng 3 (DLQ): lỗi vĩnh viễn — dữ liệu sai, bug trong code

Tầng 2 hay bị bỏ qua, nhưng nó giải quyết một trường hợp rất thực tế: message tới trước dữ liệu nó cần.

Dịch vụ A publish LeadConverted
Dịch vụ B nhận và cố nạp Lead -> chưa có, vì đồng bộ chậm

Chỉ có tầng 1: retry 3 lần trong 37 giây -> vẫn chưa có -> DLQ
Có tầng 2: thử lại sau 5 phút -> dữ liệu đã có -> THÀNH CÔNG

Ba lỗi cấu hình phổ biến:

1. Retry vô hạn:

r.Interval(int.MaxValue, TimeSpan.FromSeconds(5));     // KHÔNG BAO GIỜ làm thế này

2. Không phân biệt lỗi thoáng qua và vĩnh viễn — retry một ValidationException 3 lần rồi redelivery 3 lần nữa là 6 lần thử cho một lỗi chắc chắn không đổi (bài 17.6).

3. Không ai theo dõi DLQ — và đây là lỗi nghiêm trọng nhất.

Quy trình xử lý message trong DLQ:

_meter.CreateObservableGauge("dlq_message_count", () =>
_rabbitAdmin.LaySoMessageAsync("lead-converted_error").Result,
unit: "message");
Cảnh báo khi dlq_message_count > 0 quá 15 phút.

Ngưỡng là 0, không phải một con số khác: mỗi message trong DLQ là một sự kiện nghiệp vụ không bao giờ xảy ra, và nó cần một con người nhìn vào.

Ba bước xử lý:

# 1. Xem message và lý do
rabbitmqadmin get queue=lead-converted_error count=10 ackmode=reject_requeue_true
payload: {"EventId":"...","LeadId":"abc-999","GiaTri":5000000}
properties.headers:
MT-Fault-Message: Sequence contains no elements
MT-Fault-StackTrace: at Crm.Consumers.LeadConvertedConsumer.Consume(...)
MT-Fault-Timestamp: 2026-09-25T08:14:22Z
MT-Fault-RetryCount: 6

MassTransit lưu nguyên nhân trong header — không cần tìm trong log.

# 2. Phân loại
Dữ liệu sai (LeadId không tồn tại) -> KHÔNG replay, ghi nhận và xoá
Bug trong code -> sửa code, RỒI replay
Lỗi hạ tầng tạm thời -> replay ngay
// 3. Replay sau khi đã sửa nguyên nhân
public async Task ReplayDlqAsync(string tenQueue, int soLuong, CancellationToken ct)
{
var messages = await _rabbitAdmin.LayMessageAsync($"{tenQueue}_error", soLuong, ct);

foreach (var m in messages)
{
_logger.LogInformation("Replay message {Id} từ DLQ, lỗi gốc: {Loi}",
m.MessageId, m.Headers["MT-Fault-Message"]);

await _bus.Publish(m.Payload, m.MessageType, ct);
await _rabbitAdmin.AckAsync(m, ct);
}
}
// Endpoint quản trị, có phân quyền
app.MapPost("/admin/dlq/{queue}/replay", async (
string queue, int soLuong, IDlqService svc, CancellationToken ct) =>
{
await svc.ReplayDlqAsync(queue, soLuong, ct);
return Results.Ok();
})
.RequireAuthorization("Admin");

Hai cảnh báo khi replay:

1. Consumer PHẢI idempotent — message có thể đã xử lý một phần trước khi lỗi
2. Replay hàng loạt có thể gây đợt dội -> replay theo lô nhỏ, có độ trễ

Và quy trình phòng ngừa — viết vào runbook:

## DLQ có message

### Bước 1 — Phân loại (15 phút)
Đọc `MT-Fault-Message` của 5–10 message đầu. Cùng một lỗi hay nhiều lỗi khác nhau?

### Bước 2 — Xác định nguyên nhân
| Lỗi | Loại | Hành động |
|---|---|---|
| `Sequence contains no elements` | Dữ liệu tham chiếu không tồn tại | Kiểm tra thứ tự publish |
| `JsonException` | Hợp đồng message đổi | Xem [bài 17.11](17.11-advanced-notes.mdx) |
| `DbUpdateException` unique | Consumer chưa idempotent | Xem [bài 17.2](17.2-messaging-fears-lost-duplicates.mdx) |
| `TimeoutException` | Hạ tầng | Replay sau khi hạ tầng ổn |

### Bước 3 — Sửa nguyên nhân, rồi replay
Replay TRƯỚC khi sửa nguyên nhân = message quay lại DLQ.

### Bước 4 — Ghi nhận
Mỗi lần DLQ có message, ghi vào nhật ký sự cố. Nếu cùng một nguyên nhân
lặp lại ba lần, đó là vấn đề hệ thống cần sửa tận gốc.

Và một thiết kế tốt hơn cho trường hợp dữ liệu tham chiếu chưa tồn tại:

public async Task Consume(ConsumeContext<LeadConvertedV1> ctx)
{
var lead = await _db.Leads.FirstOrDefaultAsync(l => l.Id == ctx.Message.LeadId, ct);

if (lead is null)
{
// Không ném exception — lên lịch thử lại tường minh
if (ctx.GetRedeliveryCount() < 5)
{
_logger.LogInformation("Lead {Id} chưa tồn tại, thử lại sau 5 phút",
ctx.Message.LeadId);
await ctx.Defer(TimeSpan.FromMinutes(5));
return;
}

_logger.LogWarning("Lead {Id} không tồn tại sau 5 lần thử — bỏ qua",
ctx.Message.LeadId);
return; // ACK, không vào DLQ
}

// ...
}

Cách này phân biệt rõ hai tình huống: "chưa có, chờ thêm" và "không có, bỏ qua" — thay vì để cả hai rơi vào cùng một exception và cùng một DLQ.

Tự kiểm tra​

Câu hỏi thường gặp

Ba chỉ số nào đủ để chẩn đoán queue dồn?

Tổng số message đang chờ, số message chưa được xác nhận tức đang xử lý, và số consumer đang gắn. Ba con số này đối chiếu với bảng chẩn đoán cho ra năm nguyên nhân phổ biến.

Tổ hợp nào cho thấy thiếu consumer hoặc prefetch quá nhỏ?

Số message chờ rất cao nhưng số message đang xử lý rất thấp và số consumer ít. Nó nghĩa là consumer chỉ nhận một hoặc vài message một lúc nên không thể rút kịp.

Tổ hợp nào cho thấy poison message?

Số message chờ cao, số đang xử lý dao động, và message liên tục quay lại queue. Nhìn qua trông như consumer đang làm việc nhưng thực ra nó xử lý cùng một message mãi mãi.

Điều kiện tiên quyết trước khi tăng mức song song là gì?

Consumer phải idempotent, vì xử lý song song làm tăng khả năng hai message liên quan chạy đồng thời và va vào nhau. Ngoài ra phải kiểm tra database chịu được số truy vấn song song mới.

Vì sao cảnh báo theo ngưỡng tuyệt đối là sai?

Vì nó báo động giả vào lúc cao điểm bình thường và im lặng khi một consumer duy nhất chết trong lúc lưu lượng thấp. Cảnh báo theo xu hướng và theo tuổi message cũ nhất phản ánh đúng vấn đề hơn.

Vì sao tự co giãn theo CPU không phù hợp với consumer?

Vì consumer chờ I/O có CPU rất thấp ngay cả khi queue đang dồn hàng trăm nghìn message. Co giãn theo độ sâu queue phản ánh đúng khối lượng công việc còn lại.

Kết luận​

Ba điều đáng nhớ nhất:

  1. Ba chỉ số đủ để chẩn đoán hầu hết ca queue dồn.
  2. Queue dồn là triệu chứng của consumer, gần như không bao giờ là lỗi của broker.
  3. Cảnh báo theo tuổi message, co giãn theo độ sâu queue. Cả hai đều không phải CPU.

Tham khảo​

Điều hướng​