KDS (Kinesis Data Streams)

Kinesis Data Streams (KDS)
- Amazon Kinesis Data Streams (KDS) là dịch vụ thu thập và xử lý dữ liệu dạng streaming theo thời gian thực, với khả năng scaling linh hoạt theo số lượng shard. Đây là lựa chọn phổ biến cho các hệ thống cần ingest log, event, hoặc telemetry với độ trễ thấp (low-latency) và độ tin cậy cao.

- Thời gian lưu trữ (Retention): từ 1 ngày đến tối đa 365 ngày.
- Khả năng xử lý lại dữ liệu (Replay / Reprocess): cho phép đọc lại dữ liệu trong khoảng thời gian lưu trữ.
- Tính bất biến (Immutability): một khi dữ liệu đã được ghi vào Kinesis thì không thể xóa được.
- Tính thứ tự (Ordering): các record có cùng
partition keysẽ được ghi vào cùng một shard, đảm bảo thứ tự. - Producer (bên ghi dữ liệu):
- AWS SDK
- Kinesis Producer Library (KPL)
- Kinesis Agent
- Consumer (bên đọc dữ liệu):
- Tự phát triển (Write your own):
Kinesis Client Library (KCL)AWS SDK
- Dịch vụ được quản lý (Managed services):
AWS LambdaKinesis Data FirehoseKinesis Data Analytics
- Tự phát triển (Write your own):
Capacity Modes
Provisioned Mode
- Bạn tự chọn số lượng shard cần cấp phát, có thể scale thủ công hoặc thông qua API
- Mỗi shard có băng thông ghi (write throughput) là 1MB/s hoặc 1000 bản ghi/giây
- Mỗi shard có băng thông đọc (read throughput) là 2MB/s, áp dụng cho consumer theo mô hình classic hoặc enhanced fan-out
- Tính phí theo số lượng shard được cấp phát mỗi giờ
On-demand Mode
- Không cần cấp phát hay quản lý dung lượng
- Mặc định cấp phát với dung lượng 4MB/s ghi vào hoặc 4000 bản ghi/giây
- Tự động scale dựa trên lưu lượng truy cập cao nhất quan sát được trong 30 ngày gần nhất
- Tính phí theo stream mỗi giờ và dữ liệu vào/ra tính theo GB
Security
- Kiểm soát truy cập / phân quyền bằng IAM policies
- Mã hóa trong quá trình truyền (in-flight) sử dụng các endpoint HTTPS
- Mã hóa khi lưu trữ (at rest) sử dụng KMS (Key Management Service)
- Có thể tự implement mã hóa/giải mã phía client (phức tạp hơn, cần handle thủ công)
- VPC Endpoint khả dụng để Kinesis có thể truy cập nội bộ trong VPC
- Giám sát các API Call bằng CloudTrail

Phân tích sơ đồ
- EC2 Instance nằm trong Private Subnet và giao tiếp thông qua VPC Endpoint sử dụng HTTPS để đảm bảo bảo mật khi truyền dữ liệu (encryption in flight).
- VPC Endpoint đóng vai trò như một gateway giúp kết nối an toàn giữa private subnet và Kinesis Data Stream mà không cần đi qua internet.
- Sau khi dữ liệu được gửi tới Kinesis Stream, dữ liệu sẽ được ghi vào các shard (Shard 1, 2, 3).
- Dữ liệu trong các shard được mã hóa khi lưu trữ (encryption at rest) bằng KMS.
- Ngoài ra, AWS CloudTrail có thể được sử dụng để theo dõi và ghi log toàn bộ các API call đến Kinesis nhằm mục đích audit và bảo mật.
Sắp xếp dữ liệu vào Kinesis
- Hãy tưởng tượng bạn có 100 xe tải (truck_1, truck_2, ..., truck_100) đang di chuyển trên đường và liên tục gửi vị trí GPS của chúng vào AWS.
- Bạn muốn tiêu thụ dữ liệu theo đúng thứ tự của từng xe tải, để có thể theo dõi chuyển động của chúng một cách chính xác.
- Vậy làm thế nào để gửi dữ liệu đó vào Kinesis?
- 👉 Câu trả lời: Gửi dữ liệu bằng cách sử dụng một giá trị Partition Key là "truck_id"
Cùng một khóa (key) sẽ luôn được gửi đến cùng một shard.

Giải thích sơ đồ – Kinesis Stream với 3 Shards:
- Mỗi xe tải được đại diện bởi một màu khác nhau và một ID (ví dụ: truck 1 là màu vàng, truck 2 là màu xanh lá, v.v.)
- Khi gửi dữ liệu vào Kinesis Stream, bạn gán Partition Key = truck_id.
- AWS Kinesis sẽ sử dụng Partition Key này để quyết định shard nào sẽ nhận dữ liệu.
- Dữ liệu từ cùng một truck (ví dụ: truck_1) sẽ luôn được ghi vào cùng một shard – điều này đảm bảo thứ tự (ordering) cho từng nguồn dữ liệu riêng biệt.
- Dữ liệu từ nhiều truck khác nhau sẽ được phân phối đều trên các shard, nhưng vẫn đảm bảo dữ liệu của mỗi truck được giữ nguyên thứ tự theo thời gian.
- Nếu không dùng partition key đúng, dữ liệu có thể bị gửi đến nhiều shard → mất thứ tự giữa các records thuộc cùng 1 thực thể (truck).
Bạn thấy bài này thế nào?