</>Học Dev
Bài học

Tuần 3 - Ngày 1: Amazon Kinesis Data Streams

Tuần 3 – Ngày 1

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

PRODUCERSSTREAMCONSUMERS(app,agent,(Lambda,KCLapp,IoT,logs)Shard1Firehose,Flink)Shard2PutRecordShard3GetRecords/(partition...SubscribeToShardkeyhashRecord:partitionshard)key+data+seq#Retention:24h(default)tiđa365ngày,replayđưc
  • 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ềuGiới hạn
Write1 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

ProvisionedOn-demand
Quản lý shardTự tính và resharding (split/merge)AWS tự scale
Chi phíTheo shard-hour — rẻ hơn khi traffic ổn định, dự đoán đượcTheo 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 AgentCài trên server, tail log file gửi vào stream
AWS servicesCloudWatch 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, WriteProvisionedThroughputExceeded
  • GetRecords.IteratorAgeMilliseconds — độ tụt hậu consumer
  • ReadProvisionedThroughputExceeded — nhiều consumer chung 2MB/s → cân nhắc EFO

6. Kinesis vs SQS (phân biệt nhanh)

Kinesis Data StreamsSQS
Mô hìnhStream, nhiều consumer đọc cùng dữ liệu, replayQueue, mỗi message 1 consumer xử rồi xóa
Thứ tựTheo partition key trong shardFIFO queue (giới hạn throughput)
Retention24h-365 ngày, replayTối đa 14 ngày, không replay sau xóa
Use caseAnalytics realtime, fan-out, replayDecouple, task queue, buffering đơn

Câu hỏi ôn tập

  1. 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.

  2. Producer nhận ProvisionedThroughputExceededException dù 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.

  3. IteratorAgeMilliseconds tă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ý.

  4. Ứ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.

  5. 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.

  6. 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-recordget-records thử 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


Tiếp theo: Amazon Data Firehose