Mục tiêu học tập
- Hiểu kiến trúc Kinesis Data Streams: shard, partition key, sequence number
- Phân biệt provisioned vs on-demand capacity mode
- Nắm producer/consumer options: SDK, KPL/KCL, enhanced fan-out
- Chẩn đoán các lỗi kinh điển: hot shard, ProvisionedThroughputExceeded, iterator age
1. Kiến trúc cơ bản
- Record ≤ 1MB; được gán sequence number trong shard
- Partition key quyết định record vào shard nào (MD5 hash) — cùng key → cùng shard → giữ thứ tự theo key
- Dữ liệu lưu lại và replay được trong retention window — khác biệt cốt lõi so với Firehose
Giới hạn per shard (phải thuộc)
| Chiều | Giới hạn |
|---|---|
| Write | 1 MB/s hoặc 1,000 records/s |
| Read (shared) | 2 MB/s tổng cho mọi consumer, 5 GetRecords/s |
| Read (enhanced fan-out) | 2 MB/s cho MỖI consumer (push qua HTTP/2, ~70ms) |
2. Capacity modes
| Provisioned | On-demand | |
|---|---|---|
| Quản lý shard | Tự tính và resharding (split/merge) | AWS tự scale |
| Chi phí | Theo shard-hour — rẻ hơn khi traffic ổn định, dự đoán được | Theo GB in/out — hợp traffic gai, khó dự đoán |
| Exam keyword | "predictable traffic, cost-effective" | "unpredictable/spiky, least operational overhead" |
3. Producers
| Cách | Đặc điểm |
|---|---|
AWS SDK (PutRecord/PutRecords) | Đơn giản; PutRecords batch tới 500 records |
| KPL (Kinesis Producer Library) | Batch + aggregation (gộp nhiều record nhỏ vào 1 Kinesis record) → throughput cao, chi phí thấp; đổi lại thêm độ trễ (RecordMaxBufferedTime); retry tự động |
| Kinesis Agent | Cài trên server, tail log file gửi vào stream |
| AWS services | CloudWatch Logs subscription, EventBridge, IoT Core... |
KPL aggregation: cần KCL (hoặc deaggregation library) phía consumer để tách record. Đề nói "hàng triệu record nhỏ, tối ưu throughput" → KPL.
4. Consumers
| Cách | Đặc điểm |
|---|---|
| Lambda (event source mapping) | Poll giúp, batch size/window, parallelization factor tới 10 batch song song/shard; error handling: bisect, retry, DLQ/on-failure destination |
| KCL (Kinesis Client Library) | App tự quản (EC2/ECS); checkpoint vào DynamoDB table, tự chia shard giữa các worker (lease) |
| Enhanced fan-out (EFO) | Consumer đăng ký riêng, 2MB/s/shard dedicated, push HTTP/2 — dùng khi nhiều consumer hoặc cần latency thấp; có phí thêm |
| Firehose / Managed Flink | Đọc stream làm nguồn delivery/analytics |
5. Resharding và troubleshooting
Hot shard / hot partition
Partition key kém phân tán (vd tất cả record cùng store_id lớn) → 1 shard quá tải trong khi shard khác rảnh → ProvisionedThroughputExceededException dù tổng capacity đủ.
- Fix: partition key cardinality cao hơn (thêm hậu tố ngẫu nhiên/salt), hoặc split shard nóng, hoặc chuyển on-demand (không cứu được key quá lệch)
Iterator age cao (consumer tụt hậu)
Metric GetRecords.IteratorAgeMilliseconds tăng → consumer xử lý không kịp:
- Tăng số shard (nếu Lambda: mỗi shard 1 invocation đồng thời), tăng parallelization factor, tối ưu code, tăng memory Lambda
- Kiểm tra lỗi lặp (poison record) làm retry vô hạn — cấu hình
maximum retry attempts, bisect batch, on-failure destination (SQS/SNS)
Các metric CloudWatch cần nhớ
IncomingBytes/IncomingRecords,WriteProvisionedThroughputExceededGetRecords.IteratorAgeMilliseconds— độ tụt hậu consumerReadProvisionedThroughputExceeded— nhiều consumer chung 2MB/s → cân nhắc EFO
6. Kinesis vs SQS (phân biệt nhanh)
| Kinesis Data Streams | SQS | |
|---|---|---|
| Mô hình | Stream, nhiều consumer đọc cùng dữ liệu, replay | Queue, mỗi message 1 consumer xử rồi xóa |
| Thứ tự | Theo partition key trong shard | FIFO queue (giới hạn throughput) |
| Retention | 24h-365 ngày, replay | Tối đa 14 ngày, không replay sau xóa |
| Use case | Analytics realtime, fan-out, replay | Decouple, task queue, buffering đơn |
Câu hỏi ôn tập
-
Stream 10 shard (provisioned). Throughput ghi tối đa? Nếu 3 consumer cùng đọc kiểu shared thì mỗi consumer được bao nhiêu?
Xem đáp án
Ghi: 10 × 1 MB/s = 10 MB/s (hoặc 10,000 records/s). Đọc shared: 2 MB/s per shard chia chung cho cả 3 consumer (tổng 20 MB/s cho mọi consumer cộng lại) → mỗi consumer trung bình chỉ ~0.66 MB/s/shard và dễ bị ReadProvisionedThroughputExceeded. Cần mỗi consumer 2 MB/s/shard riêng → enhanced fan-out.
-
Producer nhận
ProvisionedThroughputExceededExceptiondù tổng traffic thấp hơn tổng capacity stream. Nguyên nhân và cách xử lý?Xem đáp án
Hot shard: partition key phân bố lệch khiến 1 shard vượt 1MB/s trong khi shard khác rảnh. Xử lý: chọn partition key cardinality cao/đều hơn (vd thêm random suffix), split shard nóng, hoặc dùng backoff-retry ở producer. Chuyển on-demand không giải quyết được nếu tất cả record dồn về một key.
-
IteratorAgeMillisecondstăng liên tục với consumer Lambda. Các cách khắc phục?Xem đáp án
Consumer tụt hậu so với producer. Cách xử lý: (1) tăng parallelization factor (tới 10 batch đồng thời/shard), (2) tăng số shard để tăng tổng concurrency, (3) tăng memory/tối ưu code Lambda, (4) tăng batch size để giảm overhead mỗi invocation, (5) kiểm tra poison record gây retry lặp — đặt max retry + on-failure destination. Nếu để iterator age vượt retention, dữ liệu sẽ mất trước khi được xử lý.
-
Ứng dụng cần thứ tự sự kiện per-device cho hàng nghìn IoT device. Thiết kế partition key thế nào?
Xem đáp án
Dùng device_id làm partition key: mọi record của một device vào cùng shard nên giữ đúng thứ tự per-device; hàng nghìn device tạo cardinality đủ cao để phân tán đều giữa các shard. Kinesis không đảm bảo thứ tự giữa các partition key khác nhau — nhưng yêu cầu chỉ là per-device nên thỏa.
-
Khi nào chọn on-demand mode thay vì provisioned?
Xem đáp án
On-demand khi traffic không dự đoán được hoặc gai (spiky) và muốn tránh vận hành resharding — trả theo GB. Provisioned khi traffic ổn định/dự đoán được — trả theo shard-hour, thường rẻ hơn ở khối lượng lớn đều đặn. Từ khóa "unpredictable workload / least operational overhead" → on-demand; "steady, cost-effective" → provisioned.
-
KPL aggregation là gì và trade-off?
Xem đáp án
KPL gộp nhiều record người dùng vào một Kinesis record ≤1MB (aggregation) + batch nhiều record vào một PutRecords call (collection) → tăng mạnh throughput, giảm chi phí per-record. Trade-off: (1) thêm latency do buffer (
RecordMaxBufferedTime), (2) consumer phải deaggregate (KCL tự làm), (3) không phù hợp yêu cầu latency cực thấp.
Bài tập thực hành
- Tạo stream on-demand, dùng CLI
aws kinesis put-recordvàget-recordsthử vòng đời record - Gắn Lambda consumer với batch size 100, quan sát metric IteratorAge trong CloudWatch
- Tính số shard cần cho 5,000 records/s, mỗi record 2KB (gợi ý: bottleneck là records/s → 5 shard)
- Đọc Kinesis Data Streams quotas
Tài liệu tham khảo chính thức
- Kinesis Data Streams — key concepts
- Enhanced fan-out consumers
- Using Lambda with Kinesis
- Resharding a stream
Tiếp theo: Amazon Data Firehose