Hôm qua bên mình gặp một ca “toang” kinh điển: một khách hàng nặng 600GB dữ liệu, một tuần phát sinh 128k đơn hàng từ sàn Tiktok. Chắc bên đó bán mỹ phẩm và đang chạy sale 8/8 nên lượng đơn tăng đột biến. Hệ quả là queue bị tồn ở Kafka nội bộ, worker xử lý không kịp, message nằm chồng chất trong hàng đợi.
Bài viết này là câu chuyện mình vừa xử lý sự cố đó: tại sao message đã vào Kafka rồi thì không thể “lôi ra một phần”, và giải pháp thực tế đã dùng là gì. Để dễ theo dõi, toàn bộ hệ thống xoay quanh 3 worker:
- ProxyWorker — nhận message từ webhook bên Tiktok, đẩy vào Kafka nội bộ.
- TiktokWorker — xử lý nghiệp vụ đồng bộ đơn hàng Tiktok.
- RePublishWorker — dùng để publish lại message vào ProxyWorker.
Mô hình ban đầu: ProxyWorker nhận từ webhook, chia theo topic name
Hệ thống của bên mình xử lý đơn hàng Tiktok theo mô hình chia 2 lớp Kafka:
- Kafka External — nơi webhook bên Tiktok đẩy message đơn hàng vào.
- Kafka Internal — Kafka nội bộ, nơi TiktokWorker thật sự đọc để xử lý.
ProxyWorker lắng nghe từ webhook bên Tiktok, rồi publish message về Kafka Internal. Khi publish, nó quyết định message rơi vào topic name nào — và đây chính là điểm mấu chốt của toàn bộ câu chuyện:
flowchart TB
TW[("webhook bên Tiktok\nKafka External")]
P["ProxyWorker\n(nhận message, chọn topic)"]
INT[("Kafka Internal")]
T1["topic order-events"]
TK["TiktokWorker\n(đồng bộ đơn hàng)"]
TW -->|"message đơn hàng"| P
P -->|"publish theo topic name"| INT
INT --> T1
T1 --> TK
Cách này hoạt động tốt khi lượng message đều đặn. Nhưng nó có một điểm mù — mọi message đều đổ vào cùng một topic, và ProxyWorker không hề biết message đó thuộc khách hàng nào, nặng hay nhẹ để tách riêng.
Sự cố: queue tồn ở Kafka Internal
Khi khách hàng 600GB bật sale 8/8, một tuần phát sinh 128k đơn hàng. Toàn bộ lượng message này đổ dồn về chung một topic, ProxyWorker là nút thắt chung — mọi khách hàng đều đi qua nó. Hệ quả:
- Message xếp hàng dài trong Kafka Internal, tốc độ phát sinh nhanh hơn tốc độ TiktokWorker xử lý.
- Lag (số message đã sản sinh nhưng chưa được xử lý) tăng vọt, tính theo trăm nghìn.
- Không chỉ khách nặng bị chậm, mà các khách hàng nhỏ khác cũng bị kẹt chung một topic — vì tất cả đang ngồi cùng một hàng đợi.
flowchart TB
TW[("webhook bên Tiktok")]
P["ProxyWorker\n(nút thắt)"]
INT[("Kafka Internal\ntopic order-events TỒN 300k")]
TK["TiktokWorker\nxử lý không kịp"]
TW -->|"128k đơn/tuần"| P
P -->|"xếp hàng dài"| INT
INT -->|"đọc chậm hơn sản sinh"| TK
Vì sao không thể “lôi message ra khỏi queue”?
Ý tưởng đầu tiên của mình là: tìm ra message của khách nặng, lôi ra khỏi queue để xử lý riêng, giải phóng đường đi cho khách còn lại. Nghe đơn giản, nhưng khi hiểu cơ chế đọc của Kafka thì mới thấy không làm vậy được.
Kafka là log, không phải queue có thể “move message”
Với RabbitMQ, có tính năng move message sẵn: message nằm trong queue, mình có thể chuyển nó sang queue khác hoặc xóa đi, ngay cả khi chưa có consumer nào đọc.
Kafka thì không có tính năng này. Kafka là một log append-only: message chỉ được ghi nối tiếp vào cuối partition, không sửa, không xóa, không di chuyển theo ý muốn. Các tool trên mạng có hỗ trợ “move” kiểu này, nhưng kiểm nghiệm cần thời gian — không thể mò dùng ngay giữa lúc đang có sự cố.
Muốn message “mất” thì phải đọc hết và commit
Cách duy nhất để một message rời khỏi phạm vi xử lý của consumer group là đọc nó và commit offset. Cơ chế hoạt động như thế này:
flowchart TB
subgraph TOPIC["Topic order-events"]
M1["msg 1"]
M2["msg 2"]
M3["msg 3"]
M4["msg 4"]
M5["msg 5"]
end
C["Consumer group\n(TiktokWorker)"]
M1 -->|"đã commit ✓"| C
M2 -->|"đã commit ✓"| C
M3 -->|"đã commit ✓"| C
M4 -->|"đang đọc"| C
M5 -->|"chưa đọc (lag)"| C
- Mỗi consumer group giữ một offset — con trỏ đánh dấu “đã xử lý đến message nào”.
- Message nào trước offset đã commit thì coi như xong; message sau offset là còn tồn (lag).
- Muốn một message “biến mất” khỏi queue với consumer group đó, bắt buộc phải đọc qua nó rồi commit — con trỏ nhảy lên, message cũ không bao giờ được đọc lại nữa.
Nghĩa là với một topic đã tồn 300k message: không có cách nào “chỉ lôi message của khách nặng ra”. Muốn dọn, phải đọc toàn bộ — và nếu chỉ đọc rồi vứt đi thì phí dữ liệu, mà nếu đọc để xử lý lại thì phải đảm bảo hệ thống không bị nghẽn tiếp.
Chia 100 partition theo shop_id vẫn không tách được khách nặng
Trước sự cố, bên mình tưởng đã làm kỹ: topic order-events chia 100 partition, và ProxyWorker publish với key là shop_id — Kafka sẽ hash shop_id để chọn partition. Nhưng thực tế vẫn bị lag, vì có ba điểm mù:
- Nhiều
shop_idrơi vào chung một partition. Kafka chọn partition theohash(shop_id) % 100— với hàng trăm nghìn shop thì chắc chắn nhiều shop dùng chung một partition. Đây không phải lỗi cấu hình, mà là bản chất của hashing. - Trong một partition, mọi message được đọc tuần tự. Consumer của partition đó phải đọc lần lượt từng message. Khách nặng xếp liền một dãy message trong partition → các shop nhỏ “vô tội” ở chung partition phải chờ hết lượt của khách nặng mới tới lượt mình.
- Số consumer tối đa bằng số partition. Chia 100 partition nghĩa là tối đa 100 luồng xử lý song song. Một partition nghẽn thì cả luồng đó đứng im — 99 luồng kia có rảnh cũng không giúp được message nằm trong partition đang nghẽn.
flowchart TB
S1["shop A (nhẹ)"]
S2["shop B (nhẹ)"]
S3["shop HEAVY\n128k đơn"]
PART[("partition #7\nnhiều shop chung 1 partition")]
CON["Consumer\nđọc tuần tự → lag"]
S1 -->|"hash(shopA) % 100"| PART
S2 -->|"hash(shopB) % 100"| PART
S3 -->|"hash(heavy) % 100"| PART
PART -->|"phải chờ HEAVY"| CON
Chia partition theo shop_id chỉ giúp dàn tải giữa các partition, chứ không tách biệt được một shop nặng khỏi các shop nhẹ. Partition vẫn là “hàng đợi chung” của những shop nằm trong nó. Muốn tách hẳn, cách duy nhất là tách tenant nặng ra một topic name riêng — đó chính là lý do giải pháp bên dưới chia theo topic name thay vì tiếp tục đánh vào partition.
Giải pháp: RePublishWorker bắn ngược queue về ProxyWorker, phân tải lại từ đầu
Cách xử lý là không “move message” trong Kafka, mà đưa toàn bộ message về lại ProxyWorker, rồi để ProxyWorker phân tải lại từ đầu — lần này với quy tắc publish mới: tenant nào có config riêng thì tách hẳn ra một topic name riêng, còn lại vẫn vào topic chung. Toàn bộ quy trình gồm 4 bước:
Bước 1: Tạo RePublishWorker, dùng lại consumer group id có sẵn
Tạo RePublishWorker nhưng dùng lại đúng consumer group id mà TiktokWorker đang dùng. Lý do phải dùng lại group cũ:
- Consumer group đang giữ vị trí offset của đống message đang tồn.
- Nếu tạo group mới, nó sẽ đọc từ một mốc khác (thường là cuối topic), không đụng được đống message cũ.
Đồng thời tắt hết TiktokWorker đang xử lý business đi — để không còn ai đọc tranh queue, tránh hai bên giằng co offset. RePublishWorker trở thành người duy nhất “cầm” group này.
flowchart TB
INT[("Kafka Internal\ntopic tồn 300k")]
OLD["TiktokWorker\n(tắt hết)"]
REP["RePublishWorker\n(dùng lại group id cũ)"]
INT -->|"chỉ RePublishWorker đọc"| REP
OLD -.->|"off"| X["×"]
Bước 2: Đọc toàn bộ queue và commit
RePublishWorker đọc hết toàn bộ message từ Kafka Internal. Quan trọng nhất: cứ đọc xong message nào thì commit offset cho message đó, để nó chính thức “mất” khỏi queue với consumer group cũ. Nếu không commit, khi bật lại TiktokWorker, nó sẽ đọc lại toàn bộ từ đầu — sự cố quay về vạch xuất phát.
flowchart TB
INT[("Kafka Internal\ntopic tồn")]
REP["RePublishWorker"]
INT -->|"đọc từng message"| REP
REP -->|"commit offset\n→ message mất"| INT
Bước 3: Publish lại message vào ProxyWorker
Trong lúc đọc, RePublishWorker publish từng message vào ProxyWorker — đúng nơi ProxyWorker lắng nghe để phân tải. Payload được giữ nguyên, không sửa, chỉ chuyển từ “queue đang tồn” sang “điểm xuất phát ban đầu”:
flowchart TB
INT[("Kafka Internal\ntopic tồn")]
REP["RePublishWorker"]
P["ProxyWorker"]
INT -->|"đọc hết + commit"| REP
REP -->|"publish lại message"| P
Sau bước này, queue nội bộ sạch: toàn bộ message đã được đọc, commit, và “trở về” chỗ ProxyWorker như thể chúng vừa mới được webhook Tiktok đẩy lên.
Bước 4: Restart, ProxyWorker phân tải lại từ đầu — config riêng thì tách topic
Giờ là lúc restart service. ProxyWorker nhận lại toàn bộ message và bắt đầu phân tải lại từ đầu — nhưng lần này có config publish mới:
- Tenant nào có config riêng (khách nặng) → message rơi vào topic riêng của tenant đó.
- Các tenant còn lại → vẫn vào topic chung
order-eventsnhư bình thường, không bị kéo theo.
flowchart TB
P["ProxyWorker\n(chia theo tenant → topic)"]
INT[("Kafka Internal")]
T0["topic order-events\n(tenant thường)"]
T1["topic order-events.ten-heavy-600gb\n(tenant có config riêng)"]
TK["TiktokWorker\n(đồng bộ đơn hàng)"]
P -->|"publish theo config tenant"| INT
INT -->|"tenant thường"| T0
INT -->|"tenant có config riêng"| T1
T0 --> TK
T1 --> TK
Với cách này, khách nặng có một topic đi riêng: nó phát sinh 128k đơn/ngày cũng chỉ nghẽn topic của nó, không kéo theo khách nhỏ nào. Chỉ tenant nào thật sự cần mới bị tách topic, còn lại vẫn dùng chung — không phải “mỗi tenant một topic” tốn công quản lý vô ích.
Kết quả: queue tách ra, rồi tự chia đều lại
Sau khi hoàn tất quy trình, queue của các khách hàng còn lại và queue của “ông” 600GB dữ liệu được tách hẳn ra — không còn ai kéo ai. Và điều thú vị: lúc sau nó tự chia đều lại.
Nguyên nhân đơn giản: khi khách nặng đã có topic riêng, nó được TiktokWorker xử lý xong lượng tồn của mình, còn các tenant nhẹ thì gần như không có backlog. Tải toàn hệ thống trở về mức cân bằng — worker nào cũng nhàn, lag về 0, mọi thứ tự ổn định lại mà không cần can thiệp thêm.
Bài học
- Kafka là log, không phải queue có thể “move message”. Message đã vào topic rồi thì không lôi ra, không chuyển chỗ — muốn nó mất phải đọc hết và commit.
- Khi queue đã tồn, đừng xử lý tiếp trong đó. Phương án khả thi là dùng RePublishWorker đọc toàn bộ + commit, publish lại message vào ProxyWorker, rồi restart với quy tắc publish mới.
- Chia chung một topic là cạm bẫy khi có “big tenant”. Một khách nặng sẽ bóp nghẹt toàn bộ hàng đợi chung — phải tách tenant nặng ra một topic name riêng, cho khách nặng một làn đi riêng; còn lại vẫn dùng topic chung cho đỡ phức tạp.
- Luôn có sẵn phương án xử lý queue tồn từ trước, đừng để đến lúc “toang” mới đi tìm tool. Và đúng là cái nghèo nó giới hạn logic code — gặp mấy khách kiểu này mà không có sẵn kịch bản thì đúng là toang :v