Kafka không biết gì về nội dung tin nhắn — với nó, tất cả chỉ là mảng byte. Điều đó khiến việc đổi lược đồ dữ liệu trở thành trách nhiệm của bạn, và bài này đo cái giá khi làm sai.

Tin hỏng khoá cứng partition, cách chữa, và ba kiểu tương thích

Một tin hỏng, tám triệu ngoại lệ

Dựng một topic có 11 tin: 5 số nguyên, 1 chuỗi lạc, rồi 5 số nguyên nữa. Consumer dùng IntegerDeserializer và chạy trong 12 giây:

đọc offset 0 = 1
đọc offset 1 = 2   ...   offset 4 = 5
LỖI: RecordDeserializationException tại offset 5

KẾT QUẢ:  đọc được 5 tin
          gặp lỗi 8.306.744 lần trong 12 giây
          offset đã chốt: 5

8.306.744 ngoại lệ trong 12 giây — 692.228 mỗi giây.

Chuyện xảy ra: poll() ném ngoại lệ trước khi trả về bất cứ bản ghi nào. Vòng lặp bắt ngoại lệ rồi gọi poll() lại, gặp đúng tin đó, ném lại. Không có độ trễ nào giữa các lần, nên nó quay ở tốc độ tối đa CPU cho phép.

Ba hậu quả cùng lúc:

  • Tin 6 đến 10 không bao giờ được đọc. Partition đứng vĩnh viễn ở offset 5.
  • Consumer đốt 100% một nhân CPU không làm gì có ích.
  • Log sinh ra khổng lồ. Nếu mỗi ngoại lệ được ghi một dòng, đó là hàng gigabyte mỗi phút — và đầy đĩa của máy chạy consumer là sự cố thứ hai chồng lên sự cố thứ nhất.

Khởi động lại consumer không giúp gì: offset đã chốt là 5, nó quay lại đúng chỗ đó.

Đây là kiểu sự cố mà một thay đổi lược đồ tưởng vô hại gây ra. Ai đó đổi kiểu một trường, hoặc đổi bộ tuần tự hoá, và một tin duy nhất đủ để dừng cả đường ống — không phải làm chậm, mà dừng hẳn.

Chữa bằng bảy dòng

RecordDeserializationException mang theo đúng thông tin cần thiết: partition nào và offset nào.

catch (RecordDeserializationException e) {
    TopicPartition tp = e.topicPartition();
    long off = e.offset();
    ghiLaiTinHong(tp, off);                                     // đưa sang topic tin chết
    c.seek(tp, off + 1);                                        // nhảy qua
    c.commitSync(Map.of(tp, new OffsetAndMetadata(off + 1)));   // chốt để khỏi quay lại
}

Chạy lại trên cùng topic:

đọc offset 0..4    ->  1 2 3 4 5
BỎ QUA tin hỏng ở offset 5
đọc offset 6..10   ->  6 7 8 9 10
KẾT QUẢ: đọc được 10 tin, bỏ qua 1 tin hỏng

Ba chi tiết quan trọng trong bảy dòng đó:

Phải commitSync ngay, không đợi tới cuối vòng lặp. Nếu consumer chết giữa seek và lần chốt tiếp theo, nó sẽ quay lại đúng tin hỏng và mọi thứ lặp lại.

Phải ghi lại tin hỏng ở đâu đó — thường là một topic riêng gọi là "tin chết". Bỏ qua trong im lặng nghĩa là mất dữ liệu mà không ai biết, và bạn cũng không còn cách nào điều tra xem chuyện gì đã xảy ra.

Phải đếm số tin bị bỏ qua và cảnh báo khi vượt ngưỡng. Một tin hỏng là chuyện thường; một nghìn tin hỏng trong một phút nghĩa là producer vừa được triển khai bản sai, và bạn muốn biết ngay chứ không phải sau khi đã bỏ qua nửa triệu tin.

Trong Spring Kafka, cơ chế tương đương có sẵn: ErrorHandlingDeserializer bọc bộ giải mã thật, và DeadLetterPublishingRecoverer đẩy tin hỏng sang topic khác. Nếu bạn dùng Spring, hãy bật nó — nhưng hãy hiểu nó đang chữa đúng chuyện gì.

Ba kiểu tương thích

Khi lược đồ phải đổi, câu hỏi thật là: nâng cấp bên nào trước?

BACKWARD — consumer mới đọc được dữ liệu . Được phép: xoá trường, thêm trường có giá trị mặc định. Nâng cấp consumer trước, rồi mới tới producer. Đây là mặc định của Confluent Schema Registry và là lựa chọn đúng cho hầu hết trường hợp.

FORWARD — consumer đọc được dữ liệu mới. Được phép: thêm trường, xoá trường có giá trị mặc định. Nâng cấp producer trước. Hợp lý khi bạn kiểm soát producer nhưng không kiểm soát hết consumer — ví dụ nhiều đội khác nhau đang đọc.

FULL — cả hai chiều. Chỉ được thêm hoặc xoá trường có giá trị mặc định. Nâng cấp theo thứ tự nào cũng được. Chặt nhất, và cũng an toàn nhất.

Ba thay đổi luôn phá vỡ mọi kiểu tương thích, và chúng là nguồn của gần hết các sự cố lược đồ:

  • Đổi kiểu một trườngint thành string.
  • Đổi tên một trường — với Avro thì đây là "xoá trường cũ, thêm trường mới", và dữ liệu cũ mất giá trị.
  • Thêm trường bắt buộc không có mặc định — dữ liệu cũ không có trường đó và không có gì để điền vào.

Đổi tên là chỗ dễ mắc nhất vì nó trông vô hại. customerId thành customer_id là một thay đổi phá vỡ hoàn toàn.

Vì sao cần một cơ quan đăng ký lược đồ

Không có nó, lược đồ tồn tại dưới dạng thoả thuận miệng giữa các đội, và không có gì ngăn producer đẩy dữ liệu mà consumer không đọc nổi. Phép đo ở đầu bài là chuyện xảy ra khi thoả thuận đó bị vi phạm.

Schema Registry giữ lược đồ ở một nơi và từ chối lược đồ mới nếu nó không tương thích theo quy tắc bạn đặt. Việc kiểm tra xảy ra lúc producer khởi động, không phải lúc consumer gặp tin hỏng — chuyển sự cố từ lúc chạy sang lúc triển khai.

Mỗi tin khi đó mang thêm 5 byte đầu: một byte đánh dấu và bốn byte id lược đồ. Consumer đọc id đó, hỏi registry lấy lược đồ tương ứng (có bộ nhớ đệm), rồi giải mã.

Cái giá: thêm một dịch vụ phải vận hành, và nó nằm trên đường đi của cả producer lẫn consumer khi khởi động. Với hệ thống một đội và vài dịch vụ, JSON kèm kỷ luật đặt tên có thể là đủ. Với nhiều đội đọc chung một topic, kỷ luật miệng sẽ hỏng — và nó hỏng theo cách ở đầu bài.

Ba quy tắc thực dụng

Luôn thêm trường mới với giá trị mặc định. Đây là một quy tắc và nó gánh phần lớn công việc.

Không bao giờ đổi tên trường. Thêm trường mới, ghi cả hai một thời gian, rồi bỏ trường cũ ở một lần triển khai sau.

Không bao giờ đổi kiểu trường. Cũng vậy: trường mới, ghi song song, rồi bỏ.

Cả ba đều là cùng một mẫu: thay đổi phá vỡ được chia thành hai bước tương thích, cách nhau đủ lâu để mọi consumer kịp nâng cấp.

Thử ba mươi giây

Kiểm consumer của bạn có sống sót qua một tin hỏng không:

# gửi một tin không đúng định dạng vào topic thật (dùng môi trường thử nghiệm!)
echo "DAY-KHONG-PHAI-DINH-DANG-HOP-LE" | \
  kafka-console-producer.sh --bootstrap-server kf:9092 --topic topic-thu-nghiem

Rồi xem consumer. Nếu nó đứng lại và log bắt đầu chạy như thác, bạn vừa tìm ra một partition có thể bị khoá cứng bởi một tin duy nhất — và bảy dòng ở trên là thứ cần thêm vào.

Phần sau đo Kafka Connect: đưa dữ liệu vào và ra Kafka mà không viết mã.