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

14.9 — 8. Message Bus Intro

Tóm tắt

Message bus giải một vấn đề cụ thể: service A cần báo cho service B mà không được biết B tồn tại. Khác biệt cốt lõi phải nắm trước khi viết dòng code nào là command và event: command là "hãy làm việc này", có đúng một người nhận, và người gửi quan tâm kết quả; event là "việc này đã xảy ra", có bao nhiêu người nhận cũng được, và người phát không quan tâm ai nghe. Nhầm hai thứ này dẫn tới thiết kế sai từ gốc. Và vấn đề khó nhất không phải gửi message — mà là giao dịch kép: lưu database rồi publish message là hai thao tác, và tiến trình chết giữa chúng làm mất message. Mẫu outbox là lời giải.

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

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

  • Phân biệt command và event, và đặt tên đúng.
  • Cấu hình MassTransit với RabbitMQ.
  • Xử lý giao dịch kép bằng outbox.
  • Cấu hình retry và dead-letter queue.
  • Quyết định khi nào message bus là thừa.

Nội dung bài học​

14.9.1 — Vấn đề nó giải​

// Ghép chặt — mỗi tính năng mới phải sửa lớp này
public async Task CreateAsync(CreateCustomerRequest request, CancellationToken ct)
{
var customer = await _repo.CreateAsync(request, ct);

await _email.SendWelcomeAsync(customer.Email, ct);
await _crmSync.PushAsync(customer, ct);
await _audit.LogAsync("CustomerCreated", customer.Id, ct);
await _analytics.TrackAsync("signup", customer.Id, ct);
}

Bốn vấn đề: mỗi tính năng mới phải sửa lớp này; email chậm làm request chậm; một service lỗi làm hỏng cả thao tác; và test phải mock bốn thứ.

// Ghep long — publish roi quen
public async Task CreateAsync(CreateCustomerRequest request, CancellationToken ct)
{
var customer = await _repo.CreateAsync(request, ct);

await _bus.Publish(new CustomerCreated(customer.Id, customer.Email), ct);
}

Thêm một consumer mới không sửa dòng nào ở đây. Đây là điểm mạnh thật sự — không phải hiệu năng, mà là khả năng thêm tính năng mà không đụng code cũ.

14.9.2 — Command và event​

CommandEvent
Ý nghĩa"Hãy làm việc này""Việc này đã xảy ra"
Số người nhậnĐúng mộtKhông giới hạn
Người gửi quan tâm kết quảCóKhông
Đặt tênĐộng từ mệnh lệnh: SendWelcomeEmailQuá khứ: CustomerCreated
MassTransitSend tới endpointPublish
// COMMAND — một người nhận
public record SendWelcomeEmail(int CustomerId, string Email);
await _bus.Send(new SendWelcomeEmail(id, email), ct);

// EVENT — bao nhiêu người nhận cũng được
public record CustomerCreated(int CustomerId, string Email, DateTime OccurredAt);
await _bus.Publish(new CustomerCreated(id, email, DateTime.UtcNow), ct);

Nhầm lẫn phổ biến nhất là publish một event mà thực chất là command:

// SAI — đây là command đội lốt event
public record SendEmailRequested(string To, string Subject, string Body);
await _bus.Publish(new SendEmailRequested(...));

Nếu hai consumer cùng đăng ký, email được gửi hai lần. Event phải mô tả sự kiện đã xảy ra trong domain, không phải việc bạn muốn ai đó làm.

Quy tắc đặt tên là công cụ tốt nhất: event luôn ở quá khứ. CustomerCreated, OrderShipped, PaymentFailed. Nếu tên không tự nhiên ở quá khứ, có lẽ đó là command.

Nội dung message nên chứa gì? Hai trường phái:

// 1. Chỉ id — consumer tự nạp dữ liệu
public record CustomerCreated(int CustomerId);

// 2. Đủ dữ liệu — consumer không cần gọi ngược
public record CustomerCreated(int CustomerId, string Email, string Name, DateTime OccurredAt);

Cách 1 luôn cho dữ liệu mới nhất nhưng tạo phụ thuộc ngược — consumer phải gọi về service gốc, và nếu service đó chết thì consumer cũng kẹt.

Cách 2 giữ các service độc lập thật sự, nhưng dữ liệu là bản chụp và message lớn hơn.

Thực dụng: đủ dữ liệu cho consumer phổ biến nhất, cộng id để consumer cần thêm thì tự nạp.

14.9.3 — MassTransit với RabbitMQ​

dotnet add package MassTransit
dotnet add package MassTransit.RabbitMQ
builder.Services.AddMassTransit(x =>
{
x.AddConsumers(typeof(Program).Assembly);

x.SetKebabCaseEndpointNameFormatter(); // ten queue de doc

x.UsingRabbitMq((context, cfg) =>
{
cfg.Host(builder.Configuration.GetConnectionString("RabbitMq"));

cfg.UseMessageRetry(r => r.Exponential(
retryLimit: 5,
minInterval: TimeSpan.FromSeconds(1),
maxInterval: TimeSpan.FromMinutes(5),
intervalDelta: TimeSpan.FromSeconds(5)));

cfg.UseCircuitBreaker(cb =>
{
cb.TrackingPeriod = TimeSpan.FromMinutes(1);
cb.TripThreshold = 15; // % thất bại
cb.ActiveThreshold = 10; // so message toi thieu
cb.ResetInterval = TimeSpan.FromMinutes(5);
});

cfg.ConfigureEndpoints(context);
});
});
public sealed class SendWelcomeEmailConsumer(
IEmailService email, ILogger<SendWelcomeEmailConsumer> logger)
: IConsumer<CustomerCreated>
{
public async Task Consume(ConsumeContext<CustomerCreated> context)
{
var message = context.Message;

logger.LogInformation("Gửi email chào mừng cho {CustomerId}", message.CustomerId);

await email.SendWelcomeAsync(message.Email, context.CancellationToken);
}
}

Circuit breaker đáng chú ý: khi dịch vụ email chết, nó ngừng thử sau một số lần thất bại thay vì tiếp tục đốt message vào một thứ đã hỏng. Sau ResetInterval nó thử lại.

Không có nó, một dịch vụ ngoài chết sẽ đẩy toàn bộ hàng đợi vào dead-letter trong vài phút.

MassTransit tự tạo queue và exchange dựa trên consumer — tiện lúc phát triển, nhưng ở production nên khai rõ để kiểm soát:

cfg.ReceiveEndpoint("welcome-email", e =>
{
e.PrefetchCount = 16;
e.ConcurrentMessageLimit = 8;
e.ConfigureConsumer<SendWelcomeEmailConsumer>(context);
});

14.9.4 — Giao dịch kép và outbox​

Đây là vấn đề khó nhất:

await _db.SaveChangesAsync(ct);                        // (1) lưu database
await _bus.Publish(new CustomerCreated(id), ct); // (2) publish

Tiến trình chết giữa (1) và (2): khách hàng đã tạo, nhưng không ai biết — email không gửi, dữ liệu không đồng bộ.

Đảo thứ tự cũng không cứu được: publish trước rồi database lỗi nghĩa là consumer xử lý một khách hàng không tồn tại.

Đây là giao dịch kép (dual write), và nó không giải được bằng cách sắp xếp lại thứ tự.

// MassTransit Outbox
builder.Services.AddMassTransit(x =>
{
x.AddEntityFrameworkOutbox<CrmDbContext>(o =>
{
o.UseSqlServer();
o.UseBusOutbox();
o.QueryDelay = TimeSpan.FromSeconds(1);
});
...
});
// Bay gio ca hai nam trong MOT transaction
await _db.SaveChangesAsync(ct); // lưu customer VÀ message
await _bus.Publish(new CustomerCreated(id), ct); // ghi vào bảng outbox

Message được ghi vào bảng outbox trong cùng transaction với dữ liệu nghiệp vụ. Một tiến trình nền đọc bảng đó và publish thật.

Kết quả: hoặc cả hai cùng được lưu, hoặc không cái nào. Tiến trình chết sau commit thì message vẫn nằm trong bảng và được publish khi khởi động lại.

Cái giá: thêm một bảng, thêm độ trễ bằng QueryDelay, và message có thể publish hai lần (tiến trình chết sau khi publish, trước khi đánh dấu đã gửi) — nên consumer vẫn phải idempotent (bài 14.10).

Mẫu này đã nói ở bài 11.8 và chi tiết ở Module 17.

14.9.5 — Dead-letter queue​

Message thất bại hết số lần retry đi vào _error queue:

welcome-email          -> queue chinh
welcome-email_error -> DEAD LETTER
welcome-email_skipped -> message không có consumer

Ba điều về nó:

  1. Message ở đó mãi — không ai tự xử lý.
  2. Phải có cảnh báo khi nó có message; queue lỗi đầy mà không ai biết là tình huống phổ biến nhất.
  3. Phải có cách xử lý lại sau khi sửa nguyên nhân.
// Cấu hình queue lỗi
cfg.ReceiveEndpoint("welcome-email", e =>
{
e.UseMessageRetry(r => r.Immediate(3));
e.UseDelayedRedelivery(r => r.Intervals(
TimeSpan.FromMinutes(5), TimeSpan.FromMinutes(15), TimeSpan.FromHours(1)));

e.ConfigureConsumer<SendWelcomeEmailConsumer>(context);
});

Hai tầng retry là mẫu chuẩn: retry ngay cho lỗi thoáng qua (mạng chập chờn), redelivery có độ trễ cho lỗi lâu hơn (dịch vụ đang khởi động lại). Chỉ khi cả hai thất bại, message mới vào dead-letter.

Queue _skipped chứa message không có consumer nào — thường nghĩa là bạn đổi tên kiểu message mà quên deploy consumer, hoặc namespace không khớp.

14.9.6 — Khi nào message bus là thừa​

Tình huốngNên dùng
Một ứng dụng, cần chạy việc nềnHangfire hoặc BackgroundService
Nhiều service cần biết về một sự kiệnMessage bus
Cần đảm bảo giao hàng qua ranh giới serviceMessage bus
Gọi service khác và cần kết quả ngayHTTP
Hàng đợi trong tiến trìnhChannel — bài 6.9

Message bus thêm rất nhiều độ phức tạp: một hạ tầng phải vận hành, gỡ lỗi khó hơn nhiều (luồng xử lý không còn tuyến tính), và test tích hợp phức tạp.

Với một monolith, nó gần như luôn là thừa. Hangfire xử lý việc nền; Channel xử lý hàng đợi trong tiến trình; MediatR xử lý tách lớp trong cùng tiến trình mà không cần hạ tầng nào.

Message bus đáng khi bạn có nhiều service triển khai độc lập và cần chúng giao tiếp mà không ghép chặt. Trước mốc đó, nó là chi phí không có lợi ích.

Một lựa chọn trung gian: MassTransit chạy với in-memory transport cho môi trường dev và test, rồi đổi sang RabbitMQ khi thật sự cần — code consumer không đổi.

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

Danh sách rà soát message bus

  • •Đã phân biệt rõ command và event; event đặt tên ở thì quá khứ.
  • •Không publish event mà thực chất là command.
  • •Nội dung message đủ cho consumer phổ biến nhất, kèm id để nạp thêm.
  • •Có outbox nếu message không được phép mất.
  • •Consumer idempotent, vì outbox và retry đều có thể gây trùng.
  • •Có retry hai tầng: ngay lập tức và redelivery có độ trễ.
  • •Có circuit breaker để không đốt cả hàng đợi khi dịch vụ ngoài chết.
  • •Có cảnh báo khi dead-letter queue có message.
  • •Có quy trình xử lý lại message trong dead-letter.
  • •Đã kiểm tra queue _skipped không có message do lệch tên kiểu.
  • •Đã cân nhắc Hangfire hoặc Channel trước khi thêm message bus.

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

Bài 1 — Mất message vì không có outbox​

Viết luồng lưu database rồi publish không có outbox, kill tiến trình giữa hai bước, và xác nhận dữ liệu tồn tại nhưng consumer không bao giờ chạy. Thêm outbox và kiểm chứng.

Tiêu chí hoàn thành: bạn giải thích được vì sao đảo thứ tự hai bước không giải quyết được, và vì sao outbox thì có.

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

Gợi ý. Có hai hệ thống lưu trữ. Bạn ghi vào cả hai. Chuyện gì xảy ra nếu chết giữa chừng?

Lời giải — bản có lỗi:

public async Task TaoLeadAsync(TaoLeadRequest req, CancellationToken ct)
{
var lead = new Lead { Name = req.Name, Email = req.Email };
_db.Leads.Add(lead);
await _db.SaveChangesAsync(ct); // bước 1: database

// Tiến trình bị kill ở ĐÂY
await _bus.Publish(new LeadDaTao(lead.Id, lead.Email), ct); // bước 2: message bus
}
curl -X POST http://localhost:5000/leads -d '{"name":"Công ty ABC","email":"an@abc.com"}'
# kill ngay sau khi SaveChanges xong
kill -9 $(pgrep -f Crm.Api)
SELECT Id, Name FROM Leads ORDER BY Id DESC;
-- 1024 Công ty ABC <- lead tồn tại
[Consumer] Không nhận được message nào
-> Email chào mừng không bao giờ được gửi
-> Bản ghi trong CRM không bao giờ được đồng bộ sang hệ thống khác

Dữ liệu tồn tại nhưng mọi hệ quả của nó thì không. Và không có lỗi nào được ghi lại — từ góc nhìn của mọi log, request đã thành công.

Vì sao đảo thứ tự không giải quyết được:

await _bus.Publish(new LeadDaTao(...), ct);     // publish TRƯỚC
// kill ở đây
await _db.SaveChangesAsync(ct);
Consumer chạy, cố nạp lead 1024
-> không tồn tại trong database
-> consumer thất bại, retry, vẫn không tồn tại, vào dead-letter queue
-> hoặc tệ hơn: consumer tạo dữ liệu phái sinh cho một lead KHÔNG BAO GIỜ tồn tại

Bạn chỉ đổi hình dạng của lỗi, không loại bỏ nó:

Thứ tựNếu chết giữa chừng
Database trướcDữ liệu có, message mất — hệ quả không xảy ra
Message trướcMessage có, dữ liệu mất — consumer xử lý thứ không tồn tại

Lý do căn bản: database và message broker là hai hệ thống lưu trữ riêng biệt, không có transaction chung. Không có thứ tự nào của hai thao tác độc lập làm chúng trở thành nguyên tử.

Đây cũng chính là bài toán ở bài 14.6 — hai tướng quân, chỉ khác bối cảnh.

Và distributed transaction (2PC) không phải câu trả lời:

MSDTC / XA transaction:
- RabbitMQ, Kafka, Redis không hỗ trợ
- Khoá tài nguyên trên cả hai hệ thống trong suốt giao thức
- Coordinator chết giữa chừng -> transaction treo, phải gỡ bằng tay
- Không dùng được qua ranh giới cloud

Ngành đã rời bỏ 2PC cho loại bài toán này, và outbox là thứ thay thế nó.

Bản sửa — outbox pattern:

public class TinNhanOutbox
{
public Guid Id { get; set; }
public string LoaiMessage { get; set; } = null!;
public string NoiDung { get; set; } = null!;
public DateTime TaoLuc { get; set; }
public DateTime? GuiLuc { get; set; }
public int SoLanThu { get; set; }
}
public async Task TaoLeadAsync(TaoLeadRequest req, CancellationToken ct)
{
await using var tx = await _db.Database.BeginTransactionAsync(ct);

var lead = new Lead { Name = req.Name, Email = req.Email };
_db.Leads.Add(lead);

// Message được ghi vào CÙNG database, trong CÙNG transaction
_db.Outbox.Add(new TinNhanOutbox
{
Id = Guid.NewGuid(),
LoaiMessage = nameof(LeadDaTao),
NoiDung = JsonSerializer.Serialize(new LeadDaTao(lead.Id, lead.Email)),
TaoLuc = DateTime.UtcNow,
});

await _db.SaveChangesAsync(ct);
await tx.CommitAsync(ct); // một điểm commit duy nhất
}
public class OutboxDispatcher : BackgroundService
{
protected override async Task ExecuteAsync(CancellationToken ct)
{
while (!ct.IsCancellationRequested)
{
var chuaGui = await _db.Outbox
.Where(m => m.GuiLuc == null && m.SoLanThu < 10)
.OrderBy(m => m.TaoLuc)
.Take(100)
.ToListAsync(ct);

foreach (var m in chuaGui)
{
try
{
await _bus.Publish(GiaiMa(m), ct);
m.GuiLuc = DateTime.UtcNow;
}
catch (Exception ex)
{
m.SoLanThu++;
_logger.LogWarning(ex, "Không gửi được message {Id}, lần {Lan}", m.Id, m.SoLanThu);
}
}

await _db.SaveChangesAsync(ct);
await Task.Delay(TimeSpan.FromSeconds(1), ct);
}
}
}

Vì sao outbox giải quyết được. Vì nó biến hai thao tác trên hai hệ thống thành một thao tác trên một hệ thống:

Trước:   [ghi database]  +  [gửi message]     hai hệ thống, không nguyên tử
Sau: [ghi database + ghi outbox] MỘT hệ thống, MỘT transaction
rồi sau đó, tách biệt:
[đọc outbox -> gửi message] có thể thất bại và thử lại tuỳ thích

Transaction của database bảo đảm: hoặc cả lead lẫn message đều được ghi, hoặc không cái nào. Không còn trạng thái dở dang.

Khoảnh khắc "chết giữa chừng" vẫn tồn tại — nhưng giờ nó nằm ở một chỗ vô hại:

Dispatcher gửi message thành công, rồi chết trước khi đánh dấu GuiLuc
-> khởi động lại, thấy message chưa đánh dấu, gửi LẠI
-> message trùng lặp, KHÔNG phải message mất

Outbox đổi mất message thành message trùng — và message trùng thì xử lý được bằng consumer idempotent, còn message mất thì không xử lý được bằng gì cả.

Bảo đảm của outbox: AT-LEAST-ONCE
Nghĩa vụ kèm theo: consumer PHẢI idempotent

Trong thực tế, dùng outbox có sẵn của MassTransit:

services.AddMassTransit(x =>
{
x.AddEntityFrameworkOutbox<CrmDbContext>(o =>
{
o.UseSqlServer();
o.UseBusOutbox();
o.QueryDelay = TimeSpan.FromSeconds(1);
o.DuplicateDetectionWindow = TimeSpan.FromMinutes(30);
});
});
// Publish trông như bình thường, nhưng đi vào outbox thay vì thẳng ra broker
await _publishEndpoint.Publish(new LeadDaTao(lead.Id, lead.Email), ct);
await _db.SaveChangesAsync(ct); // cả hai được commit cùng nhau

Nó cài sẵn bốn thứ mà bản tự viết ở trên còn thiếu: khoá để nhiều instance không cùng gửi một message, dọn message cũ, theo dõi thứ tự, và inbox phía consumer để khử trùng lặp.

Ba chi phí của outbox, nên biết trước:

  1. Độ trễ. Message không được gửi ngay mà chờ dispatcher quét — thường 100 ms tới vài giây. Với luồng cần phản hồi tức thì, đây là chi phí thật.
  2. Tải ghi lên database. Mỗi message là một INSERT cộng một UPDATE. Với hệ thống nhiều message, bảng outbox có thể thành điểm nóng.
  3. Bảng phải được dọn. Không dọn thì nó lớn vô hạn, và truy vấn WHERE GuiLuc IS NULL chậm dần vì phải quét ngày càng nhiều dòng đã gửi.
DELETE FROM Outbox WHERE GuiLuc IS NOT NULL AND GuiLuc < DATEADD(DAY, -7, SYSUTCDATETIME());

Và một index để truy vấn của dispatcher luôn nhanh bất kể bảng lớn thế nào:

CREATE INDEX IX_Outbox_ChuaGui ON Outbox(TaoLuc) WHERE GuiLuc IS NULL;

Index có lọc là lựa chọn đúng ở đây: nó chỉ chứa các dòng chưa gửi, nên kích thước của nó tỉ lệ với số message đang chờ, không phải với kích thước bảng.


Bài 2 — Event bị dùng như command​

Publish một message kiểu SendEmailRequested với hai consumer đăng ký, và xác nhận email gửi hai lần.

Tiêu chí hoàn thành: bạn phân biệt được command và event bằng ba tiêu chí, và biết đặt tên message theo đúng loại.

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

Gợi ý. Publish gửi tới bao nhiêu consumer? Send thì sao?

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

public record SendEmailRequested(string Email, string TieuDe, string NoiDung);
public class EmailConsumer : IConsumer<SendEmailRequested>
{
public async Task Consume(ConsumeContext<SendEmailRequested> ctx)
=> await _sender.GuiAsync(ctx.Message.Email, ctx.Message.TieuDe, ctx.Message.NoiDung);
}

// Người khác thêm consumer này sáu tháng sau, để lưu bản sao email
public class LuuTruEmailConsumer : IConsumer<SendEmailRequested>
{
public async Task Consume(ConsumeContext<SendEmailRequested> ctx)
{
await _luuTru.GhiAsync(ctx.Message);
await _sender.GuiAsync(ctx.Message.Email, ctx.Message.TieuDe, ctx.Message.NoiDung);
// Người viết tưởng chỉ mình consumer này xử lý message đó
}
}
await _publishEndpoint.Publish(new SendEmailRequested("an@abc.com", "Chào", "..."), ct);
[EmailConsumer]     Đã gửi email tới an@abc.com
[LuuTruEmailConsumer] Đã gửi email tới an@abc.com <- LẦN THỨ HAI

Publish gửi tới mọi consumer đã đăng ký cho kiểu message đó. Với event, đó là hành vi đúng và mong muốn. Với một yêu cầu hành động, đó là thực hiện hành động nhiều lần.

Ba tiêu chí phân biệt command và event:

CommandEvent
1. Ý nghĩa"Hãy làm X" — yêu cầu"X đã xảy ra" — thông báo
2. Số người nhậnĐúng mộtKhông hoặc nhiều
3. Người gửi biết gìBiết ai sẽ xử lý và mong nó xảy raKhông biết và không quan tâm

Tiêu chí 3 là tiêu chí quyết định, và nó là câu hỏi dễ trả lời nhất:

"Nếu không ai xử lý message này, có phải là lỗi không?"

CÓ -> command. Người gửi phụ thuộc vào việc nó được xử lý.
KHÔNG -> event. Người gửi chỉ thông báo; ai quan tâm thì nghe.

Áp vào ví dụ: nếu SendEmailRequested không được ai xử lý thì email không được gửi — đó là lỗi. Vậy nó là command, và nó đang bị publish như event.

Bản sửa:

// Command — gửi tới MỘT endpoint xác định
public record GuiEmail(string Email, string TieuDe, string NoiDung);

var endpoint = await _sendEndpointProvider.GetSendEndpoint(new Uri("queue:gui-email"));
await endpoint.Send(new GuiEmail("an@abc.com", "Chào", "..."), ct);
[EmailConsumer] Đã gửi email tới an@abc.com

Dù có bao nhiêu consumer đăng ký cho kiểu GuiEmail, chỉ một endpoint nhận được, và trong endpoint đó chỉ một consumer xử lý.

Cách đặt tên theo loại message — tên phải làm cho việc dùng sai trở nên khó chịu khi đọc:

LoạiQuy ướcVí dụ
CommandĐộng từ ở thể mệnh lệnhGuiEmail, TaoLead, HuyDonHang, TruKho
EventDanh từ + động từ quá khứLeadDaTao, DonHangDaHuy, EmailDaGui, ThanhToanDaNhan
await _sendEndpoint.Send(new GuiEmail(...));       // "gửi email" — đọc như một mệnh lệnh
await _publishEndpoint.Publish(new EmailDaGui(...)); // "email đã gửi" — đọc như một sự kiện

Cái tên SendEmailRequested ở đầu bài nằm lưng chừng: nó mô tả "một yêu cầu gửi email đã được đưa ra", nghe như event, trong khi ý định là command. Chính sự mơ hồ đó dẫn tới việc dùng Publish.

Tách hẳn về namespace để không thể nhầm:

namespace Crm.Contracts.Commands { public record GuiEmail(...); }
namespace Crm.Contracts.Events { public record EmailDaGui(...); }

Cấu hình MassTransit theo đó:

cfg.ReceiveEndpoint("gui-email", e =>
{
e.ConfigureConsumer<EmailConsumer>(context); // command: một endpoint, một consumer
});

cfg.ReceiveEndpoint("luu-tru-email", e =>
{
e.ConfigureConsumer<LuuTruEmailConsumer>(context); // nghe event EmailDaGui
});

Thiết kế lại luồng cho đúng:

// 1. Command: yêu cầu gửi
await _sendEndpoint.Send(new GuiEmail(email, tieuDe, noiDung), ct);

// 2. Consumer thực hiện, rồi PHÁT event
public class EmailConsumer : IConsumer<GuiEmail>
{
public async Task Consume(ConsumeContext<GuiEmail> ctx)
{
await _sender.GuiAsync(ctx.Message.Email, ctx.Message.TieuDe, ctx.Message.NoiDung);
await ctx.Publish(new EmailDaGui(ctx.Message.Email, DateTime.UtcNow));
}
}

// 3. Ai quan tâm thì nghe event — thêm bao nhiêu cũng được
public class LuuTruEmailConsumer : IConsumer<EmailDaGui> { }
public class ThongKeEmailConsumer : IConsumer<EmailDaGui> { }

Giờ thêm consumer mới là việc an toàn: không consumer nào có thể vô tình gửi email lần nữa, vì EmailDaGui không phải một yêu cầu.

Đây là điểm quan trọng nhất của bài. Khác biệt thật giữa command và event không nằm ở cú pháp Send hay Publish — nó nằm ở hướng của sự phụ thuộc:

Command:  người gửi phụ thuộc vào người nhận
"Tôi cần việc này được làm."
Thêm người nhận thứ hai = làm việc đó hai lần = LỖI

Event: người nhận phụ thuộc vào người gửi
"Tôi thông báo chuyện này đã xảy ra."
Thêm người nhận thứ hai = thêm một phản ứng = BÌNH THƯỜNG

Hệ quả cho kiến trúc: event cho phép mở rộng hệ thống mà không sửa code hiện có, command thì không. Đó là lý do trong một hệ thống hướng sự kiện, phần lớn message nên là event, và command chỉ nên xuất hiện ở ranh giới nơi một thành phần thật sự yêu cầu một thành phần khác làm việc cụ thể.

Một test bảo vệ:

[Fact]
public async Task Command_chi_duoc_xu_ly_mot_lan()
{
await _harness.Start();
await _harness.Bus.Send(new GuiEmail("an@abc.com", "Chào", "..."));

(await _harness.Consumed.Any<GuiEmail>()).Should().BeTrue();
_harness.Consumed.Select<GuiEmail>().Count().Should().Be(1);
}

Bài 3 — Circuit breaker cho consumer​

Làm dịch vụ email luôn thất bại, bắn 100 message, và đếm số lần consumer thử gọi có và không có circuit breaker.

Tiêu chí hoàn thành: bạn giải thích được vì sao retry đơn thuần làm sự cố nặng hơn, và chọn được tham số circuit breaker có cơ sở.

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

Gợi ý. Dịch vụ email đang quá tải. Retry gửi thêm cái gì tới nó?

Lời giải — không có circuit breaker:

cfg.ReceiveEndpoint("gui-email", e =>
{
e.UseMessageRetry(r => r.Interval(3, TimeSpan.FromSeconds(5)));
e.ConfigureConsumer<EmailConsumer>(context);
});
100 message × 4 lần gọi (1 lần đầu + 3 retry) = 400 lời gọi tới dịch vụ email
Mỗi lời gọi timeout sau 30 giây
Tổng thời gian: 400 × 30 giây / 8 consumer song song = 25 phút
Dịch vụ email: nhận 400 request trong lúc đang quá tải

Có circuit breaker:

cfg.ReceiveEndpoint("gui-email", e =>
{
e.UseCircuitBreaker(cb =>
{
cb.TrackingPeriod = TimeSpan.FromMinutes(1);
cb.TripThreshold = 15; // ngắt khi 15% thất bại
cb.ActiveThreshold = 10; // cần ít nhất 10 lời gọi mới tính
cb.ResetInterval = TimeSpan.FromMinutes(5);
});

e.UseMessageRetry(r => r.Interval(3, TimeSpan.FromSeconds(5)));
e.ConfigureConsumer<EmailConsumer>(context);
});
10 lời gọi đầu: thất bại thật, mỗi lần 30 giây timeout
Mạch NGẮT
90 message còn lại: thất bại NGAY, không gọi dịch vụ email
Tổng: 10 lời gọi thay vì 400

Giảm 40 lần số lời gọi, và thời gian từ 25 phút xuống dưới một phút.

Vì sao retry đơn thuần làm sự cố nặng hơn. Vì trong phần lớn sự cố, nguyên nhân là quá tải, và retry là thêm tải:

t=0     Dịch vụ email quá tải, bắt đầu timeout
t=0 Consumer retry -> thêm request vào một hệ thống đang quá tải
t=30s Retry cũng timeout -> retry lần nữa
t=60s Tải thực tế = tải bình thường × 4 (do retry)
-> dịch vụ email không có cơ hội nào để hồi phục
t=5m Kết nối cạn ở cả hai phía

Đây là vòng phản hồi dương: càng nhiều lỗi thì càng nhiều retry, càng nhiều retry thì càng nhiều lỗi. Nó không tự dừng cho tới khi một bên sụp hẳn.

Và có một hiệu ứng thứ hai, ít người để ý: consumer bị chiếm chỗ. Tám consumer song song, mỗi cái chờ timeout 30 giây, nghĩa là hàng đợi email không xử lý được gì khác trong suốt thời gian đó — kể cả những message lẽ ra thành công.

Circuit breaker cắt vòng lặp ở chỗ đúng: khi tỷ lệ lỗi vượt ngưỡng, ngừng gọi hẳn trong một khoảng, cho hệ thống kia không gian để hồi phục.

Đóng (Closed):   gọi bình thường, đếm tỷ lệ lỗi
Ngắt (Open): KHÔNG gọi, thất bại ngay lập tức
Thử (Half-open): sau ResetInterval, cho một số ít lời gọi đi qua
- thành công -> đóng lại
- thất bại -> ngắt tiếp

Chọn tham số có cơ sở, không chọn bừa:

ActiveThreshold — số lời gọi tối thiểu trước khi tính tỷ lệ. Không có nó, một lỗi đơn lẻ lúc hệ thống mới khởi động (1 lỗi / 1 lời gọi = 100%) sẽ ngắt mạch ngay.

Đặt bằng khoảng lưu lượng của 10–30 giây ở mức bình thường.
Hàng đợi 5 message/giây -> ActiveThreshold khoảng 50–150.

TripThreshold — phần trăm lỗi để ngắt. Phụ thuộc vào tỷ lệ lỗi bình thường của dịch vụ đó:

Dịch vụ rất ổn định (lỗi nền < 1%):     ngưỡng 10–15%
Dịch vụ có lỗi nền vài phần trăm: ngưỡng 25–30%
Dịch vụ hay lỗi lặt vặt (API bên thứ ba): ngưỡng 50%

Quy tắc: ngưỡng phải cao hơn hẳn tỷ lệ lỗi bình thường, nếu không mạch ngắt ngẫu nhiên trong lúc mọi thứ vẫn ổn. Hãy đo tỷ lệ lỗi nền trước khi đặt số.

TrackingPeriod — cửa sổ tính tỷ lệ. Ngắn quá thì nhạy với biến động nhất thời; dài quá thì phản ứng chậm.

Thường là 1–5 phút. Nên dài hơn TrackingPeriod của các retry bên trong,
để một chuỗi retry không tự nó lấp đầy cửa sổ.

ResetInterval — chờ bao lâu trước khi thử lại. Nên dài hơn thời gian hồi phục điển hình của dịch vụ đó:

Dịch vụ tự khởi động lại trong 30 giây -> ResetInterval 1 phút
Sự cố cần người can thiệp -> 5–10 phút
Dịch vụ bên thứ ba có rate limit -> theo cửa sổ rate limit của họ

Thứ tự đặt middleware là chi tiết quyết định:

e.UseCircuitBreaker(...);      // NGOÀI
e.UseMessageRetry(...); // TRONG

Đặt ngược lại thì retry chạy bên ngoài circuit breaker, và ba retry sẽ tạo ra ba lần "mạch ngắt, thất bại ngay" — không gọi dịch vụ, nhưng cũng không được lợi gì từ việc ngắt mạch, vì message vẫn bị tiêu thụ và đẩy vào dead-letter.

Với thứ tự đúng, message gặp mạch đang ngắt sẽ không bị mất: MassTransit đưa nó trở lại hàng đợi, và nó được thử lại khi mạch đóng.

Ba thành phần nên có cùng nhau:

cfg.ReceiveEndpoint("gui-email", e =>
{
// 1. Ngắt mạch khi dịch vụ đang hỏng
e.UseCircuitBreaker(cb =>
{
cb.TrackingPeriod = TimeSpan.FromMinutes(1);
cb.TripThreshold = 15;
cb.ActiveThreshold = 10;
cb.ResetInterval = TimeSpan.FromMinutes(5);
});

// 2. Giới hạn tốc độ, để không tự làm quá tải dịch vụ kia
e.UseRateLimit(100, TimeSpan.FromMinutes(1));

// 3. Retry với backoff và jitter cho lỗi thoáng qua
e.UseMessageRetry(r =>
{
r.Exponential(3,
minInterval: TimeSpan.FromSeconds(1),
maxInterval: TimeSpan.FromSeconds(30),
intervalDelta: TimeSpan.FromSeconds(5));
r.Ignore<ValidationException>(); // lỗi nghiệp vụ: retry vô nghĩa
r.Ignore<EmailKhongHopLeException>();
});

e.ConfigureConsumer<EmailConsumer>(context);
});

Hai dòng Ignore đáng chú ý: phân biệt lỗi thoáng qua và lỗi vĩnh viễn là điều kiện để retry có nghĩa. Một địa chỉ email sai định dạng sẽ sai mãi mãi — retry ba lần chỉ tốn thời gian và làm nhiễu số liệu của circuit breaker.

Nên retry:     timeout, lỗi mạng, HTTP 429, HTTP 503, deadlock database
KHÔNG retry: lỗi validate, HTTP 400, HTTP 401, HTTP 404, vi phạm ràng buộc

Và giám sát trạng thái mạch — một mạch ngắt lâu mà không ai biết là một luồng nghiệp vụ đang dừng trong im lặng:

_meter.CreateObservableGauge("circuit_breaker.state",
() => new Measurement<int>(_mach.TrangThai == TrangThaiMach.Ngat ? 1 : 0,
new KeyValuePair<string, object?>("endpoint", "gui-email")));
Cảnh báo khi mạch ở trạng thái ngắt liên tục quá 10 phút.

Tự kiểm tra​

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

Command khác event thế nào?

Command là hãy làm việc này, có đúng một người nhận, và người gửi quan tâm kết quả, đặt tên bằng động từ mệnh lệnh. Event là việc này đã xảy ra, có bao nhiêu người nhận cũng được, người phát không quan tâm ai nghe, và đặt tên ở thì quá khứ.

Publish một command dưới dạng event gây gì?

Nếu hai consumer cùng đăng ký thì việc đó được làm hai lần, ví dụ email gửi hai lần. Quy tắc đặt tên là công cụ phát hiện tốt nhất: nếu tên không tự nhiên ở thì quá khứ thì có lẽ đó là command.

Nội dung message nên chứa id hay đủ dữ liệu?

Chỉ id thì luôn có dữ liệu mới nhất nhưng tạo phụ thuộc ngược, consumer phải gọi về service gốc và kẹt nếu service đó chết. Đủ dữ liệu thì các service độc lập thật sự nhưng là bản chụp. Thực dụng là đủ cho consumer phổ biến nhất, kèm id để consumer khác tự nạp thêm.

Giao dịch kép là gì và outbox giải thế nào?

Lưu database và publish message là hai thao tác riêng, tiến trình chết giữa chúng làm mất message, và đảo thứ tự cũng không cứu được. Outbox ghi message vào một bảng trong cùng transaction với dữ liệu nghiệp vụ, rồi một tiến trình nền đọc bảng đó và publish thật.

Outbox có đảm bảo message chỉ gửi một lần không?

Không. Tiến trình có thể chết sau khi publish nhưng trước khi đánh dấu đã gửi, nên message publish hai lần. Đó là at-least-once, và consumer vẫn phải idempotent.

Khi nào message bus là thừa?

Với một monolith thì gần như luôn thừa: Hangfire xử lý việc nền, Channel xử lý hàng đợi trong tiến trình, và MediatR tách lớp trong cùng tiến trình mà không cần hạ tầng. Message bus đáng khi có nhiều service triển khai độc lập cần giao tiếp mà không ghép chặt.

Kết luận​

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

  1. Event ở thì quá khứ, command là mệnh lệnh. Nhầm hai thứ là sai từ gốc.
  2. Giao dịch kép làm mất message. Outbox là lời giải duy nhất đúng.
  3. Với monolith, message bus thường là thừa. Hangfire hoặc Channel trước đã.

Tham khảo​

Điều hướng​