Bỏ qua để đến nội dung

Integration & Messaging – SQS, SNS, Kinesis

Khi bạn bắt đầu deploy nhiều ứng dụng, chắc chắn chúng sẽ cần nói chuyện với nhau. Có hai mô hình giao tiếp giữa các ứng dụng:

  1. Synchronous (đồng bộ) – ứng dụng gọi trực tiếp ứng dụng khác. Ví dụ Buying Service gọi thẳng Shipping Service.
  2. Asynchronous / Event based (bất đồng bộ / dựa trên sự kiện) – ứng dụng gửi vào queue, ứng dụng kia lấy ra từ queue. Ví dụ Buying ServiceQueueShipping Service.

Giao tiếp đồng bộ trở nên rất có vấn đề khi có đột biến lưu lượng (spike). Hãy tưởng tượng: bình thường bạn cần encode 10 video, đột nhiên phải encode 1000 video cùng lúc. Nếu Buying Service gọi trực tiếp, dịch vụ phía sau sẽ sập hoặc timeout.

Trong trường hợp đó, tốt hơn là decouple (tách rời) các ứng dụng:

  • Dùng SQS: mô hình queue.
  • Dùng SNS: mô hình pub/sub.
  • Dùng Kinesis: mô hình real-time streaming.

Điểm mấu chốt: các dịch vụ này scale độc lập với ứng dụng của bạn!

Queue (hàng đợi) hoạt động rất trực quan: nhiều Producer gửi message vào SQS Queue, nhiều Consumer đi poll (hỏi lấy) message từ queue đó về xử lý.

Amazon SQS Standard Queue là dịch vụ lâu đời nhất của AWS (hơn 10 năm), là dịch vụ fully managed, dùng để decouple ứng dụng. Các thuộc tính cần nhớ:

  • Throughput không giới hạn, số lượng message trong queue không giới hạn.
  • Retention mặc định của message: 4 ngày, tối đa 14 ngày.
  • Độ trễ thấp (dưới 10 ms khi publish và khi receive).
  • Giới hạn 256 KB mỗi message gửi đi.
  • Có thể có message trùng lặp (at least once delivery — thỉnh thoảng xảy ra).
  • Có thể có message sai thứ tự (best effort ordering).
  • Gửi vào SQS bằng SDK (API SendMessage).
  • Message được lưu bền trong SQS cho tới khi consumer xóa nó.
  • Message retention: mặc định 4 ngày, tối đa 14 ngày.
  • Message có thể tối đa 256 KB.
  • Ví dụ: gửi một đơn hàng để xử lý, gồm Order id, Customer id, và bất kỳ attribute nào bạn muốn.
  • SQS standard: throughput không giới hạn.
  • Consumer có thể chạy trên EC2 instance, server thường, hoặc AWS Lambda.
  • Poll SQS để lấy message — nhận tối đa 10 message mỗi lần.
  • Xử lý message (ví dụ insert message vào một RDS database).
  • Xóa message bằng API DeleteMessage.

Khi có nhiều EC2 instance cùng làm consumer:

  • Consumer nhận và xử lý message song song.
  • At least once delivery.
  • Best-effort message ordering.
  • Consumer xóa message sau khi xử lý xong.
  • Bạn có thể scale consumer theo chiều ngang để tăng throughput xử lý.

Một queue dài nghĩa là công việc đang dồn lại — đó chính là tín hiệu tốt nhất để scale số lượng consumer. Cách làm chuẩn:

  • SQS Queue phát ra CloudWatch Metric – Queue Length, tên metric là ApproximateNumberOfMessages.
  • Một CloudWatch Alarm được đặt trên metric đó; khi vượt ngưỡng (breach), alarm được kích hoạt.
  • Alarm ra lệnh scale cho Auto Scaling Group chứa các EC2 instance đang poll message.

Đây là mẫu kiến trúc kinh điển: front-end web app (nằm trong một Auto Scaling Group) nhận request từ người dùng và gọi SendMessage vào SQS Queue (queue này scale vô hạn). Back-end processing application (cũng trong một Auto Scaling Group riêng) gọi ReceiveMessages để lấy message ra xử lý. Nhờ vậy hai tầng scale hoàn toàn độc lập, và front-end không bị chậm theo back-end.

Nếu ứng dụng ghi trực tiếp vào database (Amazon RDS, Amazon Aurora, Amazon DynamoDB) thì khi tải quá lớn, một số transaction có thể bị mất.

Giải pháp: đặt SQS Queue làm buffer. Ứng dụng phía trước enqueue message (SendMessage), tầng phía sau dequeue message (ReceiveMessages) rồi mới insert vào database với tốc độ mà database chịu được. Queue hấp thụ toàn bộ đột biến, không transaction nào bị mất.

  • Encryption:
    • Mã hóa đường truyền (in-flight) bằng HTTPS API.
    • Mã hóa at-rest bằng KMS keys.
    • Client-side encryption nếu client muốn tự mã hóa/giải mã.
  • Access Controls: IAM policies để điều chỉnh quyền truy cập SQS API.
  • SQS Access Policies (tương tự S3 bucket policy):
    • Hữu ích cho truy cập cross-account vào SQS queue.
    • Hữu ích để cho phép các dịch vụ khác (SNS, S3…) ghi vào SQS queue.

Đây là một trong những khái niệm bị hỏi nhiều nhất về SQS. Sau khi một message được consumer poll, nó trở nên vô hình với các consumer khác — nếu không như vậy thì hai consumer sẽ cùng xử lý một message.

  • Mặc định, message visibility timeout là 30 giây.
  • Nghĩa là message có 30 giây để được xử lý.
  • Sau khi visibility timeout kết thúc, message “hiện lại” trong SQS.

Nhìn theo thời gian: request ReceiveMessage đầu tiên → message được trả về và bắt đầu visibility timeout → trong khoảng đó các ReceiveMessage khác không nhận được message đó → hết timeout, ReceiveMessage tiếp theo nhận lại message một lần nữa.

Các hệ quả cần nhớ:

  • Nếu message không được xử lý xong trong visibility timeout, nó sẽ bị xử lý hai lần.
  • Consumer có thể gọi API ChangeMessageVisibility để xin thêm thời gian.
  • Nếu visibility timeout quá cao (hàng giờ) và consumer bị crash, việc xử lý lại sẽ mất rất lâu.
  • Nếu visibility timeout quá thấp (vài giây), ta có thể nhận message trùng lặp.

Khi consumer hỏi queue mà queue đang rỗng, nó có thể “đợi” message tới thay vì trả về rỗng ngay rồi lại hỏi tiếp. Việc này gọi là Long Polling.

  • Long Polling giảm số lượng API call tới SQS, đồng thời tăng hiệu quả và giảm độ trễ của ứng dụng.
  • Thời gian đợi có thể từ 1 giây tới 20 giây (20 giây là lựa chọn nên dùng).
  • Long Polling tốt hơn Short Polling.
  • Có thể bật Long Polling ở mức queue hoặc ở mức API bằng tham số WaitTimeSeconds.

FIFO = First In First Out — message được giữ đúng thứ tự trong queue: producer gửi 1 2 3 4 thì consumer nhận đúng 1 2 3 4.

  • Throughput bị giới hạn: 300 msg/s khi không batching, 3000 msg/s khi có batching.
  • Khả năng gửi chính xác một lần (exactly-once send), bằng cách loại bỏ trùng lặp dựa trên Deduplication ID.
  • Message được consumer xử lý theo thứ tự.
  • Thứ tự được đảm bảo theo Message Group ID — mọi message trong cùng group thì được sắp thứ tự; đây là tham số bắt buộc.

Nếu bạn cần gửi một message tới rất nhiều người nhận thì sao? Cách “direct integration” là Buying Service gọi lần lượt từng thứ: email notification, Fraud Service, Shipping Service, một SQS Queue… — càng thêm người nhận thì code của Buying Service càng rối.

Amazon SNS đảo ngược việc đó thành Pub/Sub: Buying Service chỉ publish vào một SNS Topic, và các bên quan tâm thì subscribe vào topic đó.

  • “Event producer” chỉ gửi message tới một SNS topic duy nhất.
  • bao nhiêu “event receiver” (subscription) lắng nghe topic đó cũng được.
  • Mỗi subscriber của topic nhận được tất cả message (ghi chú: có tính năng mới cho phép lọc message).
  • Tối đa 12.500.000 subscription mỗi topic.
  • Giới hạn 100.000 topic.

Các loại subscriber: SQS, Lambda, Kinesis Data Firehose, HTTP(S) Endpoints, SMS & Mobile Notifications, Emails.

SNS tích hợp với rất nhiều dịch vụ AWS

Phần tiêu đề “SNS tích hợp với rất nhiều dịch vụ AWS”

Nhiều dịch vụ AWS có thể gửi dữ liệu trực tiếp vào SNS để làm thông báo: CloudWatch Alarms, S3 Bucket (Events), Auto Scaling Group (Notifications), CloudFormation (State Changes), AWS Budgets, Lambda, AWS DMS (New Replica), DynamoDB, RDS Events

  • Topic Publish (dùng SDK):
    1. Tạo một topic.
    2. Tạo một (hoặc nhiều) subscription.
    3. Publish vào topic.
  • Direct Publish (cho SDK của mobile app):
    1. Tạo một platform application.
    2. Tạo một platform endpoint.
    3. Publish vào platform endpoint.
    • Hoạt động với Google GCM, Apple APNS, Amazon ADM

Phần bảo mật của SNS gần như song song với SQS:

  • Encryption: mã hóa đường truyền bằng HTTPS API; mã hóa at-rest bằng KMS keys; client-side encryption nếu client muốn tự làm.
  • Access Controls: IAM policies điều chỉnh quyền truy cập SNS API.
  • SNS Access Policies (tương tự S3 bucket policy): hữu ích cho cross-account access vào SNS topic, và để cho phép các dịch vụ khác (S3…) ghi vào SNS topic.

Fan Out là mẫu kiến trúc quan trọng nhất của chương này: push một lần vào SNS, nhận được ở tất cả các SQS queue đang subscribe.

  • Hoàn toàn decoupled, không mất dữ liệu.
  • SQS mang lại: lưu bền dữ liệu (data persistence), xử lý trễ (delayed processing)retry công việc.
  • Có thể thêm SQS subscriber mới theo thời gian mà không sửa producer.
  • Nhớ đảm bảo SQS queue access policy cho phép SNS ghi vào.
  • Cross-Region Delivery: hoạt động với SQS Queue ở region khác.

Có một giới hạn của S3 mà fan-out giải quyết rất gọn: với cùng một tổ hợp event type (ví dụ object create) và prefix (ví dụ images/), bạn chỉ có thể có MỘT S3 Event rule.

Nếu bạn muốn gửi cùng một S3 event tới nhiều SQS queue, hãy dùng fan-out: S3 gửi event tới một SNS Topic, rồi topic đó fan-out ra nhiều SQS Queue và cả Lambda Function.

Ứng dụng: SNS tới S3 qua Kinesis Data Firehose

Phần tiêu đề “Ứng dụng: SNS tới S3 qua Kinesis Data Firehose”

SNS có thể gửi tới Kinesis, nhờ đó ta có thêm một kiến trúc: Buying ServiceSNS TopicKinesis Data FirehoseAmazon S3 (hoặc bất kỳ destination nào mà KDF hỗ trợ).

FIFO = First In First Out — thứ tự message trong topic. SNS FIFO có các tính năng tương tự SQS FIFO:

  • Sắp thứ tự theo Message Group ID (mọi message trong cùng group được sắp thứ tự).
  • Deduplication bằng Deduplication ID hoặc Content Based Deduplication.
  • Có thể có cả SQS Standard và SQS FIFO queue làm subscriber.
  • Throughput bị giới hạn (bằng throughput của SQS FIFO).

Kết hợp SNS FIFO + SQS FIFO Fan Out khi bạn cần đồng thời fan out + thứ tự + loại trùng lặp.

  • Một JSON policy dùng để lọc message gửi tới các subscription của SNS topic.
  • Nếu một subscription không có filter policy, nó nhận mọi message.

Kinesis Data Streams dùng để thu thập và lưu dữ liệu streaming theo thời gian thực. Producer điển hình là Click Streams, IoT devices, Metrics & Logs, gửi qua Applications hoặc Kinesis Agent. Consumer điển hình là Lambda, Application, Amazon Data Firehose, Managed Service for Apache Flink.

Các đặc điểm:

  • Retention lên tới 365 ngày.
  • Khả năng xử lý lại (replay) dữ liệu ở phía consumer.
  • Không thể xóa dữ liệu khỏi Kinesis (cho tới khi nó hết hạn).
  • Dữ liệu tối đa 1 MB (use case điển hình là rất nhiều dữ liệu real-time “nhỏ”).
  • Đảm bảo thứ tự cho dữ liệu có cùng “Partition ID”.
  • Mã hóa at-rest bằng KMS, mã hóa đường truyền bằng HTTPS.
  • Kinesis Producer Library (KPL) để viết ứng dụng producer tối ưu.
  • Kinesis Client Library (KCL) để viết ứng dụng consumer tối ưu.
Provisioned mode On-demand mode
Cấu hình Tự chọn số shard Không cần provision hay quản lý capacity
Thông lượng Mỗi shard: 1 MB/s vào (hoặc 1000 record/giây), 2 MB/s ra Capacity mặc định 4 MB/s vào (hoặc 4000 record/giây)
Scale Scale thủ công khi tăng/giảm số shard Tự scale theo đỉnh throughput quan sát được trong 30 ngày gần nhất
Giá Trả tiền theo shard đã provision, theo giờ Trả tiền theo stream theo giờ + dữ liệu vào/ra theo GB

Amazon Data Firehose (trước đây gọi là “Kinesis Data Firehose”) là dịch vụ fully managed để nạp dữ liệu streaming vào các đích lưu trữ/phân tích.

Producer có thể là Applications, Client, Kinesis Agent, SDK, Kinesis Data Streams, Amazon CloudWatch (Logs & Events), AWS IoT, với record tối đa 1 MB. Firehose thực hiện data transformation bằng Lambda function nếu cần, rồi batch writes ra đích; dữ liệu toàn bộ hoặc chỉ phần lỗi (All or Failed data) có thể được ghi vào một S3 backup bucket.

Các đích được hỗ trợ:

  • AWS Destinations: Amazon Redshift / Amazon S3 / Amazon OpenSearch Service.
  • 3rd party: Splunk / MongoDB / Datadog / NewRelic / …
  • Custom HTTP Endpoint.

Các đặc điểm:

  • Tự động scale, serverless, trả tiền theo mức dùng.
  • Near Real-Time với khả năng buffering theo kích thước / theo thời gian.
  • Hỗ trợ CSV, JSON, Parquet, Avro, Raw Text, Binary data.
  • Chuyển đổi sang Parquet / ORC, nén bằng gzip / snappy.
  • Data transformation tùy ý bằng AWS Lambda (ví dụ CSV sang JSON).
Kinesis Data Streams Amazon Data Firehose
Thu thập dữ liệu streaming Nạp dữ liệu streaming vào S3 / Redshift / OpenSearch / 3rd party / custom HTTP
Cần viết code Producer & Consumer Fully managed
Real-time Near real-time
Provisioned / On-Demand mode Tự động scale
Lưu dữ liệu tới 365 ngày Không lưu dữ liệu
Có Replay Capability Không hỗ trợ replay

Đây là bảng so sánh mà hầu như ai thi SAA-C03 cũng nên học kỹ, vì đề bài thường mô tả hành vi rồi bắt bạn chọn dịch vụ:

SQS SNS Kinesis
Consumer “pull data” Push dữ liệu tới nhiều subscriber Standard: pull data2 MB mỗi shard
Dữ liệu bị xóa sau khi được consume Tối đa 12.500.000 subscriber Enhanced fan-out: push data2 MB mỗi shard mỗi consumer
Bao nhiêu worker (consumer) cũng được Dữ liệu không được lưu bền (mất nếu không gửi được) Có thể replay dữ liệu
Không cần provision throughput Pub/Sub Dành cho real-time big data, analytics và ETL
Đảm bảo thứ tự chỉ với FIFO queue Tối đa 100.000 topic Thứ tự ở mức shard
Có khả năng delay từng message riêng lẻ Không cần provision throughput Dữ liệu hết hạn sau X ngày
Tích hợp với SQS cho mẫu kiến trúc fan-out Provisioned mode hoặc on-demand capacity mode
FIFO capability cho SQS FIFO

SQS và SNS là các dịch vụ “cloud-native” — chúng dùng giao thức độc quyền của AWS. Trong khi đó, các ứng dụng truyền thống chạy on-premises thường dùng các giao thức mở như MQTT, AMQP, STOMP, Openwire, WSS.

Khi migrate lên cloud, thay vì phải re-engineer (viết lại) ứng dụng để dùng SQS và SNS, ta có thể dùng Amazon MQ.

  • Amazon MQ là dịch vụ managed message broker.
  • Amazon MQ không “scale” được nhiều như SQS / SNS.
  • Amazon MQ chạy trên server, có thể chạy Multi-AZ với failover.
  • Amazon MQ có cả tính năng queue (~SQS) và tính năng topic (~SNS).

Kiến trúc HA của Amazon MQ trong một region (ví dụ us-east-1): một Amazon MQ Broker ở trạng thái ACTIVE tại một Availability Zone (us-east-1a), một broker STANDBY tại AZ khác (us-east-1b), và cả hai dùng chung Amazon EFS làm storage. Khi broker active gặp sự cố, client failover sang broker standby.

Chủ đề Cần nhớ
Vì sao decouple Gọi synchronous dễ sập khi có spike; decouple bằng SQS (queue), SNS (pub/sub), Kinesis (streaming) — cả ba scale độc lập với ứng dụng
SQS Standard Throughput không giới hạn, message 256 KB, retention 4 ngày mặc định / 14 ngày tối đa, độ trễ < 10 ms, at-least-once delivery, best-effort ordering
Consume SQS Poll tối đa 10 message mỗi lần, xử lý, rồi gọi DeleteMessage; scale consumer theo chiều ngang
Visibility timeout Mặc định 30 giây; quá ngắn → trùng lặp, quá dài → xử lý lại rất chậm sau khi consumer crash; xin thêm thời gian bằng ChangeMessageVisibility
Long Polling Đợi 1–20 giây (nên dùng 20), bật ở mức queue hoặc bằng WaitTimeSeconds; ít API call hơn, độ trễ thấp hơn
SQS FIFO 300 msg/s không batching, 3000 msg/s có batching; exactly-once send qua Deduplication ID; Message Group ID bắt buộc cho thứ tự
Scale consumer CloudWatch metric ApproximateNumberOfMessages → CloudWatch Alarm → Auto Scaling Group
Bảo mật SQS HTTPS in flight, KMS at rest, client-side tùy chọn; IAM policy cộng SQS Access Policy cho cross-account và cho dịch vụ khác ghi vào
SNS Tối đa 12.500.000 subscription/topic100.000 topic; subscriber gồm SQS, Lambda, Firehose, HTTP(S), SMS/Mobile, Email; dữ liệu không được lưu bền
Publish SNS Topic publish qua SDK, hoặc direct publish tới platform endpoint cho GCM / APNS / ADM
Fan Out Push một lần vào SNS, mọi SQS queue subscribe đều nhận; nhớ mở SQS access policy cho SNS; hỗ trợ cross-region
S3 Events Chỉ một S3 event rule cho mỗi cặp event type + prefix, nên muốn tới nhiều queue thì đi qua SNS fan-out
SNS FIFO Thứ tự theo Message Group ID, dedup theo ID hoặc theo nội dung; ghép với SQS FIFO để có cả fan-out lẫn thứ tự
SNS Message Filtering JSON filter policy cho từng subscription; subscription không có policy thì nhận mọi message
Kinesis Data Streams Retention tới 365 ngày, replay được, record tối đa 1 MB, thứ tự theo Partition ID; Provisioned = 1 MB/s in (1000 rec/s)2 MB/s out mỗi shard; On-demand = 4 MB/s (4000 rec/s) mặc định, tự scale theo đỉnh 30 ngày
Amazon Data Firehose Fully managed, near real-time, không lưu dữ liệu, không replay; đích S3 / Redshift / OpenSearch / 3rd party / custom HTTP; transform bằng Lambda; chuyển Parquet/ORC, nén gzip/snappy
Amazon MQ Managed broker cho MQTT, AMQP, STOMP, Openwire, WSS; chạy trên server, Multi-AZ active/standby dùng Amazon EFS; có queue và topic nhưng không scale bằng SQS/SNS