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

17.12 — Mở rộng và đào sâu

Tóm tắt

Những chủ đề nằm ngoài phạm vi module, kèm điều kiện cụ thể để đáng dùng. Hai mục đáng cảnh báo riêng. Event sourcing là quyết định gần như không đảo ngược được: mọi thay đổi trở thành một event bất biến, nên "sửa dữ liệu sai" không còn là một câu UPDATE mà là một event bù trừ — và mọi người trong đội phải hiểu điều đó. Saga thì ngược lại: nó không phải lựa chọn mà là hệ quả bắt buộc của việc tách service — khi một quy trình nghiệp vụ đi qua nhiều service, bạn đã có saga, chỉ là nó đang ẩn trong code rải rác thay vì được mô hình hoá rõ ràng.

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

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

  • Biết saga giải bài toán gì và hai kiểu điều phối.
  • Đánh giá event sourcing có phù hợp hay không.
  • Hiểu vì sao schema registry quan trọng ở quy mô lớn.
  • Biết giới hạn thật của exactly-once trong Kafka.

Nội dung bài học​

17.12.1 — Saga​

Điều kiện kích hoạt: một quy trình nghiệp vụ ghi dữ liệu ở nhiều service.

Dat hang: Order -> Payment -> Inventory -> Shipping
Nếu Inventory THẤT BẠI, phải HOÀN LẠI Payment.
Không có transaction phân tán => cần BƯỚC BÙ TRỪ.

Hai kiểu điều phối:

ChoreographyOrchestration
Cách hoạt độngMỗi service nghe event và tự quyết địnhMột điều phối viên gọi từng bước
Xem toàn cảnh quy trìnhKhó — logic rải ở nhiều nơiDễ — một chỗ
CouplingThấpCao hơn (orchestrator biết mọi bước)
Gỡ lỗiKhóDễ hơn
Phù hợp2–3 bước đơn giản4 bước trở lên
// Orchestration voi MassTransit state machine
During(AwaitingPayment,
When(PaymentFailed)
.Publish(c => new ReleaseInventoryCommand(c.Saga.OrderId)) // bu tru
.Publish(c => new OrderCancelledEvent(c.Saga.OrderId))
.Finalize());

Điều khó nhất của saga không phải code mà là thiết kế bước bù trừ. Ba câu hỏi phải trả lời cho mỗi bước:

  • Bước này hoàn tác được không? (Đã gửi email cho khách thì không.)
  • Bù trừ thất bại thì sao? (Cần retry, và cuối cùng cần can thiệp thủ công.)
  • Trạng thái trung gian có lộ ra cho người dùng không? (Thường là có — phải thiết kế giao diện cho nó.)

Câu hỏi thứ nhất là lý do nên xếp các bước không hoàn tác được xuống cuối quy trình.

Đọc thêm: Saga pattern.

17.12.2 — Event sourcing​

Điều kiện kích hoạt: nghiệp vụ yêu cầu lịch sử thay đổi đầy đủ — kiểm toán, tài chính, pháp lý.

Thay vì lưu trạng thái hiện tại, lưu chuỗi sự kiện và tính trạng thái bằng cách phát lại:

Thay vi:  Lead { Status = "Won", Value = 5.000.000 }

Lưu: LeadCreated { Value = 3.000.000 }
ValueUpdated { Value = 5.000.000 }
LeadConverted { At = ... }
ƯuNhược
Lịch sử đầy đủ, miễn phíTruy vấn khó — cần read model riêng
Dựng lại trạng thái bất kỳ thời điểm nàoĐội phải học mô hình khác
Debug bằng cách xem chuỗi sự kiệnĐổi schema event rất khó
Audit trail đúng theo thiết kếGần như không đảo ngược được

Hai dòng cuối là lý do cần cân nhắc rất kỹ. Event là bất biến: sửa dữ liệu sai không còn là một câu UPDATE mà là thêm một event bù trừ. Mọi người trong đội — kể cả người vào sau — phải hiểu điều đó, nếu không sẽ có người "sửa nhanh" bằng cách ghi thẳng vào event store và phá vỡ tính toàn vẹn.

Trên .NET, Marten cung cấp event store trên PostgreSQL với projection tích hợp, là lựa chọn dễ tiếp cận nhất.

Lựa chọn trung gian, thường đúng hơn: giữ mô hình quan hệ bình thường và thêm bảng audit log ghi mọi thay đổi. Bạn có lịch sử đầy đủ mà không phải đổi toàn bộ mô hình lập trình. Với phần lớn CRM/ERP, đây là điểm cân bằng đúng.

17.12.3 — CQRS​

Điều kiện kích hoạt: mô hình đọc và mô hình ghi thật sự khác nhau, không chỉ khác về hình dạng DTO.

// Ghi: qua domain model, co quy tac nghiep vu
await _sender.Send(new ConvertLeadCommand(leadId), ct);

// Đọc: truy vấn thẳng, không qua domain
var summary = await _readDb.QueryAsync<LeadSummaryDto>(
"SELECT ... FROM LeadSummaryView WHERE TenantId = @tenantId", new { tenantId });

CQRS không bắt buộc phải có database riêng cho đọc. Mức nhẹ nhất — dùng chung database nhưng tách đường đọc khỏi đường ghi — đã cho phần lớn lợi ích: đường đọc không phải nạp entity, không qua change tracking, và tối ưu được độc lập (bài 16.5).

Chỉ tách database riêng khi đọc và ghi có nhu cầu scale thật sự khác nhau. Khi đó bạn nhận thêm một vấn đề: độ trễ đồng bộ giữa hai bên, và giao diện phải xử lý được việc ghi xong mà đọc lại chưa thấy.

Đọc thêm: CQRS pattern.

17.12.4 — Schema registry​

Điều kiện kích hoạt: nhiều đội cùng publish và consume event, và bạn đã gặp lần đầu tiên một breaking change lọt qua.

Schema registry lưu định nghĩa schema của mỗi loại message và từ chối publish nếu schema mới không tương thích ngược:

Producer đăng ký schema v2 -> Registry kiểm tra tương thích với v1
- Them truong tuy chon -> OK
- Xoá trường bắt buộc -> TỪ CHỐI
- Đổi kiểu dữ liệu -> TỪ CHỐI

Đây là cách tự động hoá quy tắc versioning ở bài 17.2. Với một đội nhỏ, code review đủ. Với nhiều đội, code review sẽ bỏ sót — và hậu quả là consumer của đội khác gãy lúc 2 giờ sáng.

Confluent Schema Registry phổ biến nhất trong hệ Kafka. Với RabbitMQ, thường dùng một thư viện contract dùng chung được đánh phiên bản — kém chặt hơn nhưng đơn giản hơn nhiều.

17.12.5 — Exactly-once của Kafka: giới hạn thật​

Kafka có transaction và idempotent producer, cho phép "exactly-once semantics" — nhưng phạm vi hẹp hơn tên gọi rất nhiều:

HOẠT ĐỘNG:     Kafka -> xử lý -> Kafka   (trong một transaction)
KHÔNG ÁP DỤNG: Kafka -> xử lý -> ghi SQL Server

Ngay khi consumer ghi ra một hệ thống ngoài Kafka, đảm bảo đó không còn, vì Kafka không thể tham gia transaction của database bạn.

// Bat idempotent producer — chong trung khi producer retry
var config = new ProducerConfig
{
EnableIdempotence = true, // chong trung o phia PRODUCER
Acks = Acks.All
};

EnableIdempotence đáng bật vì nó chống trùng do producer retry, nhưng nó không thay thế idempotency ở consumer. Với hệ thống .NET điển hình ghi vào SQL Server, bạn vẫn cần inbox hoặc unique constraint (bài 17.2).

17.12.6 — Thứ tự ưu tiên​

Trước khi đụng tới bất kỳ mục nào ở trên:

  1. Outbox — chống mất tin (bài 17.4).
  2. Idempotency ở consumer (bài 17.2).
  3. Retry có jitter và DLQ có người theo dõi (bài 17.6).
  4. Giám sát độ sâu queue và tuổi message (bài 17.14).

Bốn mục này giải quyết phần lớn sự cố thực tế. Saga chỉ cần khi thật sự có quy trình xuyên service; event sourcing thì hiếm khi cần.

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

Danh sách rà soát trước khi dùng mẫu nâng cao

  • •Đã có outbox và idempotency trước khi nghĩ tới saga.
  • •Nếu có quy trình xuyên service, nó được mô hình hoá rõ chứ không rải rác.
  • •Mỗi bước trong saga đều có bước bù trừ tương ứng.
  • •Bước không hoàn tác được xếp xuống cuối quy trình.
  • •Đã cân nhắc bảng audit log trước khi chọn event sourcing.
  • •Nếu dùng CQRS, đã bắt đầu từ mức nhẹ dùng chung database.
  • •Quy tắc versioning message được kiểm tra tự động hoặc trong review.
  • •Không dựa vào exactly-once của Kafka khi ghi ra database ngoài.

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

Bài 1 — Vẽ saga ẩn​

Tìm một quy trình nghiệp vụ đi qua nhiều service trong hệ thống của bạn và vẽ các bước kèm bước bù trừ tương ứng.

Tiêu chí hoàn thành: bạn vẽ được cả bước xuôi lẫn bước bù trừ, và nhận ra rằng saga đã tồn tại trong hệ thống — chỉ là nó ẩn và không đầy đủ.

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

Gợi ý. Khi một bước ở giữa thất bại, hệ thống của bạn làm gì với những bước đã xong?

Lời giải — tìm quy trình nhiều bước:

# Chuỗi event nhiều tầng thường là saga ẩn
grep -rn "INotificationHandler<\|IConsumer<" --include="*.cs" src/ \
| grep -oP "(?<=<)[A-Za-z]+(?=V?\d*>)" | sort | uniq -c | sort -rn

Quy trình "chốt lead thành khách hàng":

1. Chốt lead                    (CRM)
2. Tạo khách hàng (CRM)
3. Tạo subscription (Billing)
4. Cấp tài khoản đăng nhập (Identity)
5. Gửi email chào mừng (Notification)
6. Đồng bộ sang ERP (Integration)

Vẽ cả hai chiều:

BƯỚC XUÔI                          BƯỚC BÙ TRỪ
──────────────────────────── ──────────────────────────────
1. Chốt lead <- Đưa lead về trạng thái Qualified
2. Tạo khách hàng <- Đánh dấu khách hàng là "chưa kích hoạt"
3. Tạo subscription <- Huỷ subscription
4. Cấp tài khoản <- Vô hiệu hoá tài khoản
5. Gửi email chào mừng <- KHÔNG BÙ TRỪ ĐƯỢC <- điểm quan trọng
6. Đồng bộ ERP <- Gửi lệnh huỷ sang ERP

Bước 5 không bù trừ được, và đó là thông tin quan trọng nhất từ bài tập này: email đã gửi không thu hồi được.

Hệ quả cho thiết kế: sắp xếp các bước sao cho bước không bù trừ được nằm CUỐI CÙNG.

Thứ tự sai:  ... -> gửi email (5) -> đồng bộ ERP (6)
ERP thất bại -> phải bù trừ ngược
-> nhưng email đã gửi -> khách nhận email chào mừng
cho một tài khoản sẽ bị huỷ

Thứ tự đúng: ... -> đồng bộ ERP (5) -> gửi email (6)
ERP thất bại -> bù trừ các bước 1–4
-> email CHƯA gửi -> không ai biết gì

Saga đã tồn tại trong hệ thống của bạn — chỉ là nó ẩn:

// Đây LÀ một saga, chỉ là không ai gọi nó như vậy
public class ChotLeadHandler : INotificationHandler<LeadDaChot>
{
public async Task Handle(LeadDaChot e, CancellationToken ct)
{
await _khachHangService.TaoAsync(e.LeadId, ct); // bước 2
}
}

public class TaoKhachHangHandler : INotificationHandler<KhachHangDaTao>
{
public async Task Handle(KhachHangDaTao e, CancellationToken ct)
{
await _billingService.TaoSubscriptionAsync(e.CustomerId, ct); // bước 3
}
}
// ... và cứ thế
Đây là chuỗi event ba tầng ở [bài 16.8](../module-16-clean-architecture/16.8-in-process.mdx),
nhưng qua nhiều DỊCH VỤ — nên nó còn có thêm ba vấn đề:

1. KHÔNG CÓ trạng thái: không ai biết quy trình đang ở bước nào
2. KHÔNG CÓ bù trừ: bước 4 thất bại -> bước 1, 2, 3 vẫn nguyên
3. KHÔNG THẤY ĐƯỢC: phải đọc code của 4 dịch vụ mới biết quy trình đầy đủ

Ba câu hỏi mà không ai trả lời nhanh được:

"Lead abc-123 đã convert xong chưa?"
-> phải kiểm tra 4 database

"Vì sao khách hàng def-456 có tài khoản mà không có subscription?"
-> bước 3 thất bại lúc nào đó, không ai biết

"Bao nhiêu quy trình đang dở dang?"
-> không có cách nào biết

Sau khi làm rõ thành saga tường minh:

public class TrangThaiChotLead : SagaStateMachineInstance
{
public Guid CorrelationId { get; set; }
public string CurrentState { get; set; } = null!;

public Guid LeadId { get; set; }
public Guid? CustomerId { get; set; }
public Guid? SubscriptionId { get; set; }
public Guid? AccountId { get; set; }

public DateTime BatDauLuc { get; set; }
public DateTime? HoanTatLuc { get; set; }
public string? LyDoThatBai { get; set; }
}
public class ChotLeadSaga : MassTransitStateMachine<TrangThaiChotLead>
{
public State DangTaoKhachHang { get; private set; } = null!;
public State DangTaoSubscription { get; private set; } = null!;
public State DangCapTaiKhoan { get; private set; } = null!;
public State DangBuTru { get; private set; } = null!;
public State HoanTat { get; private set; } = null!;
public State DaHuy { get; private set; } = null!;

public ChotLeadSaga()
{
InstanceState(x => x.CurrentState);

Initially(
When(LeadDaChot)
.Then(ctx =>
{
ctx.Saga.LeadId = ctx.Message.LeadId;
ctx.Saga.BatDauLuc = DateTime.UtcNow;
})
.TransitionTo(DangTaoKhachHang)
.Send(ctx => new TaoKhachHang(ctx.Saga.CorrelationId, ctx.Saga.LeadId)));

During(DangTaoKhachHang,
When(KhachHangDaTao)
.Then(ctx => ctx.Saga.CustomerId = ctx.Message.CustomerId)
.TransitionTo(DangTaoSubscription)
.Send(ctx => new TaoSubscription(ctx.Saga.CorrelationId, ctx.Saga.CustomerId!.Value)),

When(TaoKhachHangThatBai)
.Then(ctx => ctx.Saga.LyDoThatBai = ctx.Message.LyDo)
.TransitionTo(DangBuTru)
.Send(ctx => new HuyChotLead(ctx.Saga.LeadId))); // BÙ TRỪ

During(DangTaoSubscription,
When(SubscriptionDaTao)
.Then(ctx => ctx.Saga.SubscriptionId = ctx.Message.SubscriptionId)
.TransitionTo(DangCapTaiKhoan)
.Send(ctx => new CapTaiKhoan(ctx.Saga.CorrelationId, ctx.Saga.CustomerId!.Value)),

When(TaoSubscriptionThatBai)
.Then(ctx => ctx.Saga.LyDoThatBai = ctx.Message.LyDo)
.TransitionTo(DangBuTru)
.Send(ctx => new XoaKhachHang(ctx.Saga.CustomerId!.Value)) // BÙ TRỪ bước 2
.Send(ctx => new HuyChotLead(ctx.Saga.LeadId))); // BÙ TRỪ bước 1

During(DangCapTaiKhoan,
When(TaiKhoanDaCap)
.Then(ctx => ctx.Saga.HoanTatLuc = DateTime.UtcNow)
.TransitionTo(HoanTat)
.Publish(ctx => new ChotLeadHoanTat(ctx.Saga.LeadId, ctx.Saga.CustomerId!.Value)),
// Email gửi Ở ĐÂY — bước cuối, sau khi mọi thứ đã chắc chắn

When(CapTaiKhoanThatBai)
.TransitionTo(DangBuTru)
.Send(ctx => new HuySubscription(ctx.Saga.SubscriptionId!.Value))
.Send(ctx => new XoaKhachHang(ctx.Saga.CustomerId!.Value))
.Send(ctx => new HuyChotLead(ctx.Saga.LeadId)));

During(DangBuTru,
When(BuTruHoanTat).TransitionTo(DaHuy));

// Timeout — quy trình treo quá lâu
Schedule(() => TimeoutQuyTrinh, x => x.TimeoutTokenId, s =>
{
s.Delay = TimeSpan.FromMinutes(10);
s.Received = r => r.CorrelateById(m => m.Message.CorrelationId);
});
}
}

Bốn thứ saga tường minh cho bạn mà chuỗi event không cho:

-- 1. Trạng thái hiện tại của mọi quy trình
SELECT CurrentState, COUNT(*) FROM TrangThaiChotLead GROUP BY CurrentState;
CurrentState          count
HoanTat 8421
DangTaoSubscription 12
DangBuTru 3
DaHuy 47
DangCapTaiKhoan 2
-- 2. Quy trình treo — phát hiện được
SELECT LeadId, CurrentState, BatDauLuc, LyDoThatBai
FROM TrangThaiChotLead
WHERE HoanTatLuc IS NULL AND BatDauLuc < DATEADD(MINUTE, -10, SYSUTCDATETIME());
3. Bù trừ tự động khi thất bại
4. Toàn bộ quy trình mô tả ở MỘT file, đọc được

Và cái giá của saga:

- Một khái niệm mới phải học (state machine)
- Bảng trạng thái phải lưu trữ và dọn dẹp
- Mỗi bước cần một cặp event thành công/thất bại
- Bù trừ phải được viết và TEST — và test bù trừ khó hơn test đường xuôi
- Debug phức tạp hơn

Khi nào saga xứng đáng:

XỨNG ĐÁNG:
- Quy trình qua 3+ dịch vụ
- Thất bại ở giữa để lại trạng thái không nhất quán CÓ HẬU QUẢ
- Cần biết quy trình đang ở đâu
- Quy trình kéo dài (phút, giờ, ngày)

KHÔNG XỨNG ĐÁNG:
- Quy trình 2 bước
- Thất bại ở giữa không gây hậu quả (chỉ là việc phụ không xảy ra)
- Mọi bước trong cùng một dịch vụ -> dùng transaction

Và lời khuyên quan trọng nhất: tránh cần saga bằng cách thiết kế lại ranh giới.

Nếu ba bước đầu LUÔN phải thành công cùng nhau,
có lẽ chúng thuộc về CÙNG một dịch vụ.

-> Gộp CRM và Billing thành một dịch vụ
-> Ba bước đầu thành một transaction
-> Saga chỉ còn hai bước, hoặc không cần nữa

Saga là công cụ để xử lý ranh giới dịch vụ đã có. Nếu bạn đang thiết kế một saga phức tạp cho một hệ thống chưa tách dịch vụ, hãy xem lại ranh giới trước — một transaction đơn giản hơn nhiều so với một state machine có bù trừ.


Bài 2 — Thử Marten và event sourcing​

Dựng một aggregate đơn giản với event store trên PostgreSQL và dựng lại trạng thái từ chuỗi event.

Tiêu chí hoàn thành: bạn dựng lại được trạng thái, và nêu được bốn cái giá của event sourcing cùng điều kiện để nó đáng dùng.

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

Gợi ý. Event sourcing lưu chuỗi thay đổi thay vì trạng thái hiện tại. Truy vấn "mọi lead có giá trị trên 500 triệu" thực hiện thế nào?

Lời giải:

dotnet add package Marten
// Event — bất biến, chỉ ghi thêm
public record LeadDaTao(Guid LeadId, string Ten, decimal GiaTri);
public record LeadDaGan(Guid LeadId, Guid NhanVienId);
public record LeadDaChot(Guid LeadId, DateTime ThoiDiem);
public record LeadDaHuy(Guid LeadId, string LyDo);
// Aggregate — dựng lại từ event
public class Lead
{
public Guid Id { get; private set; }
public string Ten { get; private set; } = null!;
public decimal GiaTri { get; private set; }
public Guid? NguoiPhuTrach { get; private set; }
public string TrangThai { get; private set; } = "New";
public int Version { get; private set; }

public void Apply(LeadDaTao e) { Id = e.LeadId; Ten = e.Ten; GiaTri = e.GiaTri; }
public void Apply(LeadDaGan e) { NguoiPhuTrach = e.NhanVienId; }
public void Apply(LeadDaChot e) { TrangThai = "Won"; }
public void Apply(LeadDaHuy e) { TrangThai = "Cancelled"; }
}
builder.Services.AddMarten(o =>
{
o.Connection(chuoiKetNoi);
o.Events.StreamIdentity = StreamIdentity.AsGuid;
o.Projections.Snapshot<Lead>(SnapshotLifecycle.Inline);
}).UseLightweightSessions();
// Ghi event
await using var session = _store.LightweightSession();
var leadId = Guid.CreateVersion7();

session.Events.StartStream<Lead>(leadId,
new LeadDaTao(leadId, "Công ty ABC", 5_000_000));
await session.SaveChangesAsync(ct);

session.Events.Append(leadId, new LeadDaGan(leadId, nhanVienId));
session.Events.Append(leadId, new LeadDaChot(leadId, DateTime.UtcNow));
await session.SaveChangesAsync(ct);
// Dựng lại trạng thái từ chuỗi event
var lead = await session.Events.AggregateStreamAsync<Lead>(leadId, token: ct);

Console.WriteLine($"{lead!.Ten}: {lead.TrangThai}, version {lead.Version}");
Công ty ABC: Won, version 3

Dựng lại tại một thời điểm trong quá khứ:

var leadLucTruoc = await session.Events
.AggregateStreamAsync<Lead>(leadId, version: 2, token: ct);
Công ty ABC: New, version 2        <- trước khi chốt
// Hoặc theo thời điểm
var leadHomQua = await session.Events
.AggregateStreamAsync<Lead>(leadId,
timestamp: DateTimeOffset.UtcNow.AddDays(-1), token: ct);

Đọc toàn bộ lịch sử:

var lichSu = await session.Events.FetchStreamAsync(leadId, token: ct);

foreach (var e in lichSu)
Console.WriteLine($"v{e.Version} {e.Timestamp:HH:mm:ss} {e.EventType.Name}: {e.Data}");
v1 08:14:22 LeadDaTao:  { LeadId: ..., Ten: "Công ty ABC", GiaTri: 5000000 }
v2 08:15:03 LeadDaGan: { LeadId: ..., NhanVienId: ... }
v3 09:42:17 LeadDaChot: { LeadId: ..., ThoiDiem: ... }

Đây là lợi ích lớn nhất: lịch sử đầy đủ, không mất mát, có thứ tự, và nó là NGUỒN SỰ THẬT — không phải một bảng audit có thể bị quên cập nhật.

Bốn cái giá của event sourcing:

Cái giá 1 — truy vấn khó.

-- Với mô hình thông thường
SELECT * FROM Leads WHERE GiaTri > 500000000 AND Status = 'Won';
Với event sourcing:
không có bảng Leads
-> phải dựng lại MỌI stream để lọc? -> không khả thi
-> cần PROJECTION: một bảng đọc được dựng từ event
public class LeadProjection : SingleStreamProjection<LeadDocument>
{
public void Apply(LeadDaTao e, LeadDocument d)
{
d.Id = e.LeadId; d.Ten = e.Ten; d.GiaTri = e.GiaTri; d.TrangThai = "New";
}
public void Apply(LeadDaChot e, LeadDocument d) => d.TrangThai = "Won";
}
o.Projections.Add<LeadProjection>(ProjectionLifecycle.Async);
var leads = await session.Query<LeadDocument>()
.Where(l => l.GiaTri > 500_000_000 && l.TrangThai == "Won")
.ToListAsync(ct);
-> Giờ bạn có HAI mô hình phải duy trì: event store và projection
-> Và projection là nhất quán CUỐI CÙNG với event store nếu chạy Async

Cái giá 2 — phiên bản event là vĩnh viễn.

Event đã ghi KHÔNG BAO GIỜ sửa được — đó là bản chất của event sourcing.

Đổi cấu trúc event -> phải xử lý CẢ định dạng cũ, MÃI MÃI.
// V1 đã có 2 triệu event trong store
public record LeadDaTao(Guid LeadId, string Ten, decimal GiaTri);

// V2 — cần thêm TienTe
public record LeadDaTaoV2(Guid LeadId, string Ten, decimal GiaTri, string TienTe);
o.Events.Upcast<LeadDaTao, LeadDaTaoV2>(cu =>
new LeadDaTaoV2(cu.LeadId, cu.Ten, cu.GiaTri, "VND")); // giả định cho dữ liệu cũ

Hàm upcast này phải được giữ mãi mãi, vì event cũ vẫn nằm trong store. Sau năm năm và mười lần đổi, bạn có một chuỗi upcast dài.

Cái giá 3 — dung lượng tăng vô hạn.

Mô hình thông thường:  một hàng cho mỗi lead
Event sourcing: một hàng cho MỖI THAY ĐỔI

Lead được sửa 50 lần -> 50 event
200.000 lead × trung bình 20 event = 4 triệu hàng

Và không xoá được — xoá event là phá vỡ khả năng dựng lại. Snapshot giúp tăng tốc đọc nhưng không giảm dung lượng:

o.Projections.Snapshot<Lead>(SnapshotLifecycle.Inline);

Cái giá 4 — xoá dữ liệu theo yêu cầu pháp lý rất khó.

Người dùng yêu cầu xoá dữ liệu cá nhân (GDPR)
-> event chứa tên, email, số điện thoại
-> xoá event = phá vỡ chuỗi

Giải pháp: crypto-shredding
-> mã hoá dữ liệu cá nhân trong event bằng khoá riêng cho mỗi người
-> "xoá" = huỷ khoá
-> event còn nguyên nhưng không đọc được nữa

Đây là một lớp phức tạp đáng kể phải thiết kế từ đầu — thêm vào sau rất khó.

Điều kiện để event sourcing đáng dùng — cần ÍT NHẤT hai trong bốn:

1. LỊCH SỬ là yêu cầu nghiệp vụ, không phải tiện ích
(kiểm toán tài chính, hồ sơ y tế, hợp đồng pháp lý)

2. Cần trả lời "trạng thái lúc X là gì" thường xuyên

3. Nhiều read model khác nhau từ cùng một nguồn
(CQRS với 5+ projection)

4. Bản thân nghiệp vụ đã tư duy theo sự kiện
(kế toán: bút toán; ngân hàng: giao dịch; bảo hiểm: yêu cầu bồi thường)

Và những trường hợp KHÔNG nên dùng:

- CRUD thông thường
- Chỉ cần audit log -> dùng temporal table (bài 12.9) hoặc bảng audit
- Đội chưa có kinh nghiệm -> chi phí học rất cao
- Toàn bộ hệ thống -> event sourcing cho MỌI aggregate gần như luôn sai

Cách dùng thực dụng: event sourcing cho MỘT vài aggregate cốt lõi.

Order        -> event sourcing (lịch sử quan trọng, nhiều trạng thái)
Payment -> event sourcing (yêu cầu kiểm toán)
Lead -> mô hình thông thường
Customer -> mô hình thông thường
TinhThanh -> mô hình thông thường

Marten hỗ trợ cả hai trong cùng một database — bạn dùng session.Events cho aggregate event-sourced và session.Store cho document thông thường. Đây là điểm mạnh thật của nó, và nó cho phép áp dụng dần thay vì chuyển đổi toàn bộ.

Một lựa chọn nhẹ hơn cho phần lớn trường hợp: temporal table.

ALTER TABLE Leads SET (SYSTEM_VERSIONING = ON (HISTORY_TABLE = dbo.LeadsHistory));
SELECT * FROM Leads FOR SYSTEM_TIME AS OF '2026-09-25 09:14:20' WHERE Id = 1;

Nó cho bạn "trạng thái tại một thời điểm" — lợi ích số 2 trong bốn lợi ích — mà không đổi mô hình lập trình, không cần projection, và không có ba cái giá còn lại (bài 12.9).

Với nhiều dự án nghĩ rằng mình cần event sourcing, temporal table là thứ họ thật sự cần.


Bài 3 — Kiểm tra tương thích schema​

Lấy một event đang dùng, thêm một trường bắt buộc, và xác nhận consumer cũ gãy.

Tiêu chí hoàn thành: bạn thấy consumer gãy, và biết cách để CI phát hiện thay vì production.

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

Gợi ý. Bài 17.1 đã nói về tương thích hai chiều. Bài này là về việc tự động hoá việc kiểm tra.

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

// Producer thêm trường bắt buộc
public record LeadDaChotV1(
Guid EventId, Guid LeadId, decimal GiaTri, string TienTe,
string TenantId); // BẮT BUỘC, không có mặc định
// Consumer cũ, chưa cập nhật
public record LeadDaChotV1(Guid EventId, Guid LeadId, decimal GiaTri, string TienTe);
System.Text.Json.JsonException: JSON deserialization for type 'LeadDaChotV1'
was missing required properties including: 'TenantId'.

Hoặc, nếu producer gửi trước và consumer cũ nhận: trường thừa bị bỏ qua, consumer chạy bình thường. Chiều nào gãy phụ thuộc vào ai deploy trước — và bạn không kiểm soát được điều đó.

Ba cách để CI phát hiện, theo thứ tự chi phí tăng dần:

Cách 1 — payload mẫu trong repo (rẻ nhất, hiệu quả nhất):

src/Crm.Contracts/TestData/
LeadDaChotV1/
2026-03-baseline.json
2026-06-them-tenantid.json
2026-09-them-machiendich.json
[Theory]
[MemberData(nameof(MoiPayloadMau))]
public void Kieu_hien_tai_doc_duoc_moi_payload_tung_ton_tai(string duongDan)
{
var payload = File.ReadAllText(duongDan);
var kieu = SuyRaKieuTuDuongDan(duongDan);

var act = () => JsonSerializer.Deserialize(payload, kieu, _jsonOptions);

act.Should().NotThrow(
$"payload {Path.GetFileName(duongDan)} có thể còn trong queue hoặc DLQ");
}

public static IEnumerable<object[]> MoiPayloadMau()
=> Directory.GetFiles("TestData", "*.json", SearchOption.AllDirectories)
.Select(f => new object[] { f });
Xpect: payload 2026-03-baseline.json có thể còn trong queue hoặc DLQ,
but found System.Text.Json.JsonException with message
"was missing required properties including: 'TenantId'"

Test đỏ ngay ở CI, trước khi merge.

Cách 2 — sinh payload tự động từ định nghĩa hiện tại:

[Fact]
public void Sinh_payload_mau_cho_moi_integration_event()
{
var kieuEvent = typeof(LeadDaChotV1).Assembly.GetTypes()
.Where(t => t.Namespace?.Contains("Contracts.Events") == true)
.Where(t => t.IsClass && !t.IsAbstract);

foreach (var t in kieuEvent)
{
var mau = TaoInstanceMau(t);
var json = JsonSerializer.Serialize(mau, t, _jsonOptions);

var duongDan = Path.Combine("TestData", t.Name,
$"{DateTime.UtcNow:yyyy-MM}-sinh-tu-dong.json");
Directory.CreateDirectory(Path.GetDirectoryName(duongDan)!);

if (!File.Exists(duongDan)) File.WriteAllText(duongDan, json);
}
}

Chạy test này mỗi khi thêm trường, commit file sinh ra. Sau một năm bạn có hồ sơ đầy đủ mọi định dạng — và không ai phải nhớ tạo nó bằng tay.

Cách 3 — schema registry (đắt nhất, mạnh nhất):

services.AddSingleton<ISchemaRegistryClient>(
new CachedSchemaRegistryClient(new SchemaRegistryConfig
{
Url = "http://schema-registry:8081",
}));
# Đặt chế độ tương thích cho một subject
curl -X PUT http://localhost:8081/config/lead-events-value \
-H "Content-Type: application/json" \
-d '{"compatibility": "FULL"}'
BACKWARD:  consumer mới đọc được message cũ
FORWARD: consumer cũ đọc được message mới
FULL: cả hai <- nên dùng cho hệ thống nhiều đội
NONE: không kiểm tra
# Đăng ký schema mới -> registry TỪ CHỐI nếu không tương thích
curl -X POST http://localhost:8081/subjects/lead-events-value/versions \
-H "Content-Type: application/vnd.schemaregistry.v1+json" \
-d '{"schema": "..."}'
{
"error_code": 409,
"message": "Schema being registered is incompatible with an earlier schema"
}
Ưu:    kiểm tra TẬP TRUNG, mọi dịch vụ đều phải qua
bắt được cả khi producer thuộc đội khác
Nhược: thêm một thành phần phải vận hành
chủ yếu dùng với Avro/Protobuf; với JSON thì hỗ trợ yếu hơn

So sánh ba cách:

Payload mẫuSinh tự độngSchema registry
Chi phí dựng30 phút2 giờ1–2 ngày
Thành phần thêm001
Bắt đượcĐịnh dạng đã lưuĐịnh dạng đã lưuMọi thay đổi
Bắt được producer đội khácKhôngKhôngCó
Phù hợpMột đội, JSONMột đội, nhiều eventNhiều đội, Avro/Protobuf

Với phần lớn dự án, cách 1 là đủ — và nó cho 80% lợi ích với 5% công sức.

Đưa vào CI:

- name: Kiểm tra tương thích hợp đồng message
run: dotnet test tests/Crm.Contracts.Tests --filter "Category=Compatibility"
[Trait("Category", "Compatibility")]
public class TuongThichTests { }

Và một kiểm tra bổ sung — chặn việc thêm trường bắt buộc:

[Fact]
public void Integration_event_khong_duoc_co_truong_bat_buoc_moi()
{
var kieuEvent = typeof(LeadDaChotV1).Assembly.GetTypes()
.Where(t => t.Namespace?.Contains("Contracts.Events") == true && t.IsClass);

var viPham = new List<string>();

foreach (var t in kieuEvent)
{
var ctor = t.GetConstructors().OrderByDescending(c => c.GetParameters().Length).First();

// Đếm tham số KHÔNG có giá trị mặc định
var batBuoc = ctor.GetParameters().Count(p => !p.HasDefaultValue);
var soTruongGoc = LaySoTruongGocTuBaseline(t.Name);

if (soTruongGoc is not null && batBuoc > soTruongGoc)
viPham.Add($"{t.Name}: {batBuoc} trường bắt buộc, baseline có {soTruongGoc}");
}

viPham.Should().BeEmpty(
"trường mới phải có giá trị mặc định để giữ tương thích ngược");
}

Test này bắt vấn đề tại nguồn — lúc ai đó viết string TenantId thay vì string? TenantId = null — thay vì bắt nó qua một payload mẫu bị lỗi.

Và một cảnh báo về required:

public record LeadDaChotV1
{
public required Guid EventId { get; init; }
public required string TenantId { get; init; } // required -> phá vỡ tương thích
}

required của C# 11 làm System.Text.Json ném exception khi trường thiếu — kể cả với kiểu nullable. Nó tốt cho DTO nội bộ, nhưng với integration event thì không dùng, vì nó loại bỏ khả năng bỏ qua trường thiếu.

// Cho integration event — dùng nullable với mặc định, không dùng required
public record LeadDaChotV1(
Guid EventId,
Guid LeadId,
decimal GiaTri,
string? TenantId = null);

Và xử lý trường thiếu một cách tường minh trong consumer, thay vì để serializer quyết định:

if (string.IsNullOrEmpty(e.TenantId))
{
_logger.LogWarning("Message {EventId} thiếu TenantId — message từ phiên bản cũ", e.EventId);
return; // hoặc suy ra từ dữ liệu khác
}

Tự kiểm tra​

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

Vì sao nói saga không phải lựa chọn mà là hệ quả?

Vì khi một quy trình nghiệp vụ ghi dữ liệu ở nhiều service, bạn đã có saga rồi. Câu hỏi chỉ là nó được mô hình hoá rõ ràng hay đang ẩn trong code rải rác khắp các service.

Choreography và orchestration khác nhau thế nào?

Choreography để mỗi service tự nghe event và quyết định, coupling thấp nhưng khó xem toàn cảnh và khó gỡ lỗi. Orchestration có một điều phối viên gọi từng bước, dễ xem và dễ gỡ lỗi hơn, phù hợp từ bốn bước trở lên.

Điều khó nhất của saga là gì?

Thiết kế bước bù trừ. Phải trả lời cho mỗi bước xem nó có hoàn tác được không, bù trừ thất bại thì sao, và trạng thái trung gian có lộ ra cho người dùng không. Bước không hoàn tác được nên xếp xuống cuối.

Vì sao event sourcing gần như không đảo ngược được?

Vì event là bất biến, nên sửa dữ liệu sai không còn là một câu UPDATE mà phải thêm event bù trừ. Cả đội kể cả người vào sau đều phải hiểu điều đó, nếu không sẽ có người ghi thẳng vào event store và phá vỡ tính toàn vẹn.

Lựa chọn trung gian thay cho event sourcing là gì?

Giữ mô hình quan hệ bình thường và thêm bảng audit log ghi mọi thay đổi. Bạn có lịch sử đầy đủ mà không phải đổi toàn bộ mô hình lập trình, và với phần lớn CRM thì đó là điểm cân bằng đúng.

Giới hạn thật của exactly-once trong Kafka là gì?

Nó chỉ áp dụng cho luồng Kafka sang Kafka trong một transaction. Ngay khi consumer ghi ra một database ngoài, Kafka không thể tham gia transaction đó nên đảm bảo không còn, và bạn vẫn cần inbox hoặc unique constraint.

Kết luận​

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

  1. Saga là hệ quả của việc tách service, không phải một lựa chọn thêm vào.
  2. Event sourcing gần như không đảo ngược được. Bảng audit log thường là câu trả lời đúng hơn.
  3. Exactly-once của Kafka không áp dụng khi ghi ra database ngoài.

Tham khảo​

Điều hướng​