Kafka Streams là thư viện xử lý dòng dữ liệu chạy ngay trong ứng dụng của bạn — không cần cụm riêng. Bài này đo một phép đếm đơn giản, và cả ba con số đều bất ngờ.
Đếm theo khoá
Bảy dòng cho một phép đếm hoàn chỉnh, chịu lỗi, có trạng thái:
StreamsBuilder b = new StreamsBuilder();
b.<String,String>stream("st-in")
.groupByKey()
.count()
.toStream()
.mapValues(String::valueOf)
.to("st-dem");
Chạy trên 200.000 tin với 5.000 khoá phân biệt:
đầu vào 200.000 bản ghi
đầu ra 5.340 bản ghi ít hơn 37 lần
mỗi khoá đếm được 40 đúng (200.000 / 5.000)
Kết quả đúng — mỗi khoá đếm được 40. Nhưng chỉ có 5.340 bản ghi ra, không phải 200.000.
Streams không phát một bản ghi ra cho mỗi bản ghi vào
Đây là điều tôi không mong đợi và nó thay đổi cách suy nghĩ về đầu ra.
Về mặt ngữ nghĩa, mỗi tin vào làm số đếm của một khoá thay đổi, nên "đúng" ra phải có 200.000 cập nhật. Streams có một bộ đệm bản ghi (cache.max.bytes.buffering, mặc định 10 MB) nuốt các cập nhật trung gian và chỉ phát giá trị mới nhất của mỗi khoá ở mỗi chu kỳ commit.
5.000 khoá × khoảng một lần phát mỗi khoá = 5.340. Con số khớp.
Hai hệ quả thực dụng:
Đừng dựa vào việc thấy mọi giá trị trung gian. Nếu hạ nguồn cần thấy từng bước thay đổi, phải tắt bộ đệm (cache.max.bytes.buffering=0) — và khi đó đầu ra sẽ là 200.000 bản ghi, với chi phí tương ứng.
Đây là một tính năng, không phải mất mát. Với bảng theo dõi hay bảng tổng hợp, chỉ giá trị cuối mới quan trọng, và giảm 37 lần lưu lượng ghi là món hời rất lớn.
Trạng thái nằm ở hai nơi
trên đĩa cục bộ RocksDB, 380 KB, 4 tệp SST
trong Kafka app-st-dem-KSTREAM-AGGREGATE-STATE-STORE-0000000001-changelog
cleanup.policy=compact
RocksDB trên đĩa là bản làm việc. Heap của ứng dụng chỉ 35 MB dù đang giữ 5.000 khoá — vì trạng thái nằm ngoài heap, trong RocksDB. Đây là lý do Kafka Streams giữ được trạng thái lớn hơn nhiều so với RAM mà không gây gai dọn rác.
Topic changelog trong Kafka là bản sao lưu. Streams tự tạo nó, và nó dùng cleanup.policy=compact — phần 18 đã đo: nén theo khoá khiến topic hội tụ về đúng cỡ trạng thái thay vì tăng mãi.
Cặp này là toàn bộ cơ chế chịu lỗi: máy chết thì RocksDB mất, nhưng changelog còn, và Streams dựng lại từ đó.
Khởi động lại ngay chậm hơn khởi động lại chậm — 90 lần
Tôi xoá kho trạng thái cục bộ rồi khởi động lại bốn lần liên tiếp:
| Thời gian tới RUNNING | Bản ghi khôi phục | |
|---|---|---|
| lần 1 — kho trống | 363 ms | 5.340 |
| lần 2 — ngay sau lần 1 | 44.313 ms | 0 |
| lần 3 — ngay sau lần 2 | 29.140 ms | 0 |
| lần 4 — sau khi chờ 50 giây | 463 ms | 0 |
Đọc bảng này ngược với mọi trực giác.
Dựng lại toàn bộ trạng thái từ changelog mất 363 ms. Việc mà ai cũng sợ hoá ra rất rẻ ở quy mô này.
Khởi động lại ngay lập tức mất 44 giây — và không dựng lại gì cả. Toàn bộ thời gian đó là chờ.
Nguyên nhân: Kafka Streams cố ý không rời nhóm khi đóng (internal.leave.group.on.close=false). Ý định là tốt — nó tránh một lần cân bằng lại khi ứng dụng chỉ khởi động lại nhanh. Nhưng bản mới tham gia với một danh tính thành viên khác, nên trong một khoảng thời gian nhóm có hai thành viên cho cùng số partition, và bản mới phải đợi phiên cũ hết hạn (session.timeout.ms, mặc định 45.000) mới nhận đủ partition để tới RUNNING.
Cách chữa là group.instance.id — thành viên tĩnh, phần 9 đã đo. Khi đó bản mới khai đúng danh tính cũ và nhận lại ngay phần cũ, không phải đợi gì.
group.instance.id=streams-instance-1 # lấy từ tên pod nếu chạy StatefulSet
session.timeout.ms=60000 # dài hơn thời gian khởi động ứng dụng
Nếu bạn từng thấy ứng dụng Streams "treo 45 giây mỗi lần triển khai", đây là nguyên nhân, và một dòng cấu hình xoá nó.
Gộp theo cửa sổ
s.groupByKey()
.windowedBy(TimeWindows.ofSizeWithNoGrace(Duration.ofSeconds(5)))
.count()
.toStream((w, v) -> w.key() + "@" + w.window().start())
Trên cùng dữ liệu:
5.000 bản ghi ra
khoá dạng K4952@1788029660000
kho trạng thái 440 KB (so với 380 KB khi đếm không cửa sổ)
Khoá cửa sổ mang thêm mốc bắt đầu, nên số khoá phân biệt trong kho trạng thái là số khoá × số cửa sổ đang mở. Ở đây dữ liệu được ghi trong một giây nên tất cả rơi vào một cửa sổ, và kho chỉ lớn hơn 16%.
Với dữ liệu thật trải dài, con số nhân này là thứ quyết định kích thước kho. Cửa sổ 1 phút giữ 24 giờ nghĩa là mỗi khoá có 1.440 mục — và đó là lý do retention của cửa sổ (mặc định 24 giờ) là tham số đáng nhìn trước khi kho trạng thái phình ra.
Thời gian sự kiện, không phải thời gian tường
Cửa sổ của Streams chạy theo dấu thời gian của bản ghi, không theo đồng hồ máy. Đây là điều khiến kết quả lặp lại được: xử lý lại dữ liệu cũ cho ra đúng cùng một kết quả, dù chạy vào lúc nào.
Nó cũng có nghĩa là dữ liệu đến muộn có thể rơi ra ngoài cửa sổ. ofSizeWithNoGrace là tên gọi thẳng thắn: không có thời gian ân hạn, tin đến sau khi cửa sổ đóng thì bị bỏ. ofSizeAndGrace(Duration.ofSeconds(5), Duration.ofSeconds(30)) cho thêm 30 giây ân hạn, đổi lại kho trạng thái phải giữ cửa sổ lâu hơn.
Phần sau đo chính xác chuyện gì xảy ra với dữ liệu đến muộn khi nối hai luồng.
Ba thứ Streams tự tạo mà bạn phải biết
Chạy một ứng dụng Streams là để nó tự tạo topic trong cụm của bạn:
- changelog — sao lưu kho trạng thái,
compact - repartition — khi có
groupByđổi khoá, dữ liệu phải được rải lại __consumer_offsets— như mọi consumer
Chúng mang tên bắt đầu bằng application.id, và xoá ứng dụng không xoá chúng. Công cụ kafka-streams-application-reset.sh là thứ dọn đúng cách; xoá tay dễ để sót và lần chạy sau sẽ dựng lại trạng thái từ dữ liệu cũ.
Thử ba mươi giây
Xem ứng dụng Streams của bạn đang giữ bao nhiêu trạng thái:
# kho cục bộ
du -sh /tmp/kafka-streams/<application-id>
# topic nội bộ Streams đã tạo
kafka-topics.sh --bootstrap-server kf:9092 --list | grep "^<application-id>-"
# kích thước changelog
kafka-log-dirs.sh --bootstrap-server kf:9092 --describe \
--topic-list <application-id>-...-changelog | grep -oE '"size":[0-9]+'
Nếu changelog lớn hơn nhiều so với kho cục bộ, việc nén chưa chạy tới — và min.cleanable.dirty.ratio là chỗ đầu tiên nên nhìn.
Phần sau đo phép nối hai luồng: cửa sổ 5 giây cho 0 kết quả, cửa sổ 60 giây cho 5.000 — trên cùng một tập dữ liệu.