Nối hai dòng dữ liệu là việc phổ biến nhất và cũng dễ làm sai nhất trong Kafka Streams. Bài này dựng một phép nối rồi thay đúng một con số, và kết quả đi từ đầy đủ xuống bằng không.
Phép nối
KStream<String,String> don = b.stream("j-don");
KStream<String,String> tt = b.stream("j-tt");
don.join(tt,
(x, y) -> x + "|" + y,
JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofSeconds(5)))
.to("j-out");
Nối theo khoá, và chỉ nối những cặp có dấu thời gian cách nhau không quá 5 giây.
Chạy trên 20.000 đơn hàng và 20.000 thanh toán gửi liền nhau:
20.000 kết quả D3 -> don-3|tt-3
Đúng như mong đợi.
Thay một con số, kết quả về không
Bây giờ dựng lại đúng dữ liệu đó nhưng có khoảng cách thật: gửi 5.000 đơn hàng, chờ 30 giây, rồi gửi 5.000 thanh toán cùng khoá.
| Cửa sổ nối | Kết quả |
|---|---|
| 5 giây | 0 |
| 60 giây | 5.000 |
Không có gì khác giữa hai lần chạy ngoài một tham số.
Và điều tệ nhất: trường hợp 5 giây không báo lỗi, không cảnh báo, không có dòng log nào. Topic đầu ra chỉ đơn giản là rỗng.
Đây là kiểu hỏng khó chẩn đoán nhất, vì nó trông y hệt "chưa có dữ liệu đi qua". Bạn kiểm tra topic nguồn — có dữ liệu. Kiểm tra ứng dụng — đang chạy, không lỗi. Kiểm tra topic đích — rỗng. Không có gì chỉ ra rằng 5.000 cặp vừa bị bỏ đi vì lệch nhau 30 giây.
Cửa sổ nối chạy theo thời gian sự kiện
Điểm mấu chốt: cửa sổ so dấu thời gian của bản ghi, không so thời điểm ứng dụng nhận được chúng.
Điều đó có hai mặt.
Mặt tốt: xử lý lại dữ liệu cũ cho ra đúng cùng kết quả. Chạy lại hôm nay một luồng dữ liệu của tuần trước sẽ nối đúng những cặp đã nối tuần trước — không phụ thuộc tốc độ đọc.
Mặt xấu: cửa sổ phải phản ánh khoảng cách trong nghiệp vụ, không phải khoảng cách trong hệ thống. Nếu thanh toán thường tới sau đơn hàng 2 phút, cửa sổ 30 giây sẽ không bao giờ nối được gì — dù cả hai luồng đều chạy nhanh và khoẻ.
Cửa sổ lớn hơn thì tốn gì
Câu hỏi tự nhiên tiếp theo: vậy cứ đặt cửa sổ thật dài cho chắc?
Đo trên 20.000 × 20.000 bản ghi gửi liền nhau:
| Kết quả | Heap | Kho trạng thái | |
|---|---|---|---|
| cửa sổ 5 giây | 20.000 | 53 MB | 1,5 MB |
| cửa sổ 60 giây | 20.000 | 54 MB | 1,5 MB |
Gần như không tốn thêm gì.
Lý do: kho trạng thái chỉ giữ những bản ghi thật sự đang chờ bên kia. Khi dữ liệu tới gần nhau, mỗi bản ghi được ghép ngay và không nằm lại lâu, nên độ dài cửa sổ không ảnh hưởng.
Chi phí tỉ lệ với lượng dữ liệu trong cửa sổ, không tỉ lệ với độ dài cửa sổ. Ở lưu lượng 1.000 tin/giây, cửa sổ 60 giây phải giữ 60.000 bản ghi mỗi luồng; ở lưu lượng 10 tin/giây thì chỉ 600.
Nên nguyên tắc là: cửa sổ dài không đắt nếu lưu lượng thấp, và rất đắt nếu lưu lượng cao. Cả hai con số phải cùng vào phép tính.
Phép nối giữ hai kho
join-j-out5-KSTREAM-JOINTHIS-0000000004-store-changelog
join-j-out5-KSTREAM-JOINOTHER-0000000005-store-changelog
Hai kho, một cho mỗi luồng. Mỗi bản ghi được lưu lại trong suốt cửa sổ để chờ bên kia — vì bên kia có thể tới trước hoặc sau.
Nối hai luồng lưu lượng cao với cửa sổ dài là cách nhanh nhất làm phình kho trạng thái, và cũng là cách nhanh nhất làm chậm việc dựng lại trạng thái khi khởi động lại.
Ba kiểu nối, và cái nào giữ trạng thái
| Cửa sổ | Kho trạng thái | Dùng khi | |
|---|---|---|---|
KStream ⋈ KStream |
có | hai kho, giữ theo cửa sổ | hai dòng sự kiện, tương quan theo thời gian |
KStream ⋈ KTable |
không | một kho (bảng) | làm giàu sự kiện bằng dữ liệu tra cứu |
KTable ⋈ KTable |
không | hai kho (bảng) | giữ hai bảng đồng bộ |
KStream ⋈ KTable là kiểu bị đánh giá thấp nhất và thường là kiểu đúng. Làm giàu đơn hàng bằng thông tin khách hàng không cần cửa sổ: bảng khách hàng là trạng thái hiện tại, và mỗi đơn hàng tra vào đó.
Nếu bạn đang dùng nối luồng–luồng để làm giàu dữ liệu, gần như chắc chắn bạn nên dùng luồng–bảng và bỏ hẳn bài toán cửa sổ.
Một cái bẫy của luồng–bảng: bảng phải được nạp trước khi sự kiện tới. Streams có max.task.idle.ms để bảo tác vụ chờ dữ liệu bảng, nhưng ở lần khởi động đầu tiên, sự kiện đến trước khi bảng kịp dựng sẽ tra ra null. GlobalKTable giải quyết chuyện đó — nó dựng đầy đủ trước khi xử lý bắt đầu — với cái giá là mỗi thực thể giữ toàn bộ bảng.
Chọn độ dài cửa sổ
Đừng đoán. Đo:
-- trên dữ liệu lịch sử của bạn
select percentile_cont(0.99) within group (order by
extract(epoch from thanh_toan_luc - dat_hang_luc))
from don_hang where thanh_toan_luc is not null;
Lấy p99 rồi nhân đôi. Đó là cửa sổ.
- Quá ngắn → mất kết quả trong im lặng, đúng như bảng đầu bài.
- Quá dài → kho phình theo lưu lượng, và độ trễ tăng vì nối ngoài phải chờ hết cửa sổ mới biết là không có cặp.
Và luôn thêm thời gian ân hạn cho dữ liệu đến muộn:
JoinWindows.ofTimeDifferenceAndGrace(Duration.ofMinutes(5), Duration.ofMinutes(1))
ofTimeDifferenceWithNoGrace là tên rất thẳng thắn — nó nói rõ rằng bạn đang chọn không có ân hạn. Nếu luồng của bạn có thể chậm (và mọi luồng đều có thể), một phút ân hạn là rẻ.
Thử ba mươi giây
Kiểm phép nối của bạn có đang lặng lẽ mất dữ liệu không:
# số bản ghi vào hai luồng
kafka-get-offsets.sh --bootstrap-server kf:9092 --topic luong-a
kafka-get-offsets.sh --bootstrap-server kf:9092 --topic luong-b
# số bản ghi ra
kafka-get-offsets.sh --bootstrap-server kf:9092 --topic ket-qua-noi
Nếu số ra thấp hơn nhiều so với số cặp bạn mong đợi, cửa sổ nối là chỗ đầu tiên nên nhìn — và không có gì trong log sẽ nói cho bạn biết điều đó.
Phần sau đo Kafka Connect: đưa dữ liệu vào và ra Kafka mà không viết mã.