Mục tiêu học tập
- Hiểu Amazon Managed Service for Apache Flink: kiến trúc, KPU, checkpointing
- Nắm windowing: tumbling, sliding, session — chọn đúng theo yêu cầu
- Biết các lựa chọn streaming ETL khác: Glue streaming, Lambda
- Chọn đúng công cụ xử lý stream theo độ phức tạp
1. Managed Service for Apache Flink
Tổng quan
Tên cũ Kinesis Data Analytics (đổi 8/2023). Dịch vụ chạy ứng dụng Apache Flink managed: xử lý stream stateful, latency thấp, exactly-once processing.
Đặc điểm chính
- Stateful: giữ state (đếm, session, join buffer) — được checkpoint tự động vào S3 → phục hồi exactly-once processing khi lỗi
- Scale theo KPU (Kinesis Processing Unit = 1 vCPU + 4GB); auto scaling
- Billing theo KPU-hour (chạy liên tục — không phải per-invocation như Lambda)
- Use case: windowed aggregation, anomaly detection, stream join, sessionization — logic mà Lambda đơn lẻ khó làm đúng
2. Windowing — trọng tâm thi
| Window | Định nghĩa | Ví dụ |
|---|---|---|
| Tumbling | Cửa sổ cố định, không chồng lấp | Doanh số mỗi 5 phút (00:00-00:05, 00:05-00:10) |
| Sliding (hopping) | Cửa sổ cố định, trượt/chồng lấp | "5 phút gần nhất, cập nhật mỗi 1 phút" (moving average) |
| Session | Đóng khi ngừng hoạt động quá gap | Phiên người dùng: gom event tới khi idle 15 phút |
Tumbling (5m): |----W1----|----W2----|----W3----|
Sliding (5m/1m):|----W1----|
|----W2----|
|----W3----|
Session (gap): |--W1--| gap |----W2----| gap |W3|
Từ khóa: "every N minutes, non-overlapping" → tumbling; "last N minutes updated every M" / moving average → sliding; "user session / inactivity gap" → session.
- Event time vs processing time + watermark: xử lý dữ liệu đến trễ theo thời gian sự kiện — Flink hỗ trợ đầy đủ (điểm mạnh so với xử lý thô bằng Lambda)
3. Các lựa chọn streaming ETL khác
Glue Streaming job
- Spark Structured Streaming (micro-batch, mặc định ~100s window, chỉnh được) đọc Kinesis/Kafka
- Hợp khi: đội đã dùng Glue/Spark, cần ghi vào lake (kể cả Iceberg/Hudi), chấp nhận latency giây-phút
- Cũng có checkpoint (lưu vị trí đọc) để phục hồi
Lambda
- Xử lý per-batch, stateless từ Kinesis/MSK — transform, route, alert đơn giản
- Không phù hợp: window dài, state lớn, join stream — phải tự lưu state ngoài (DynamoDB) → phức tạp, dễ sai
So sánh chọn công cụ xử lý stream
| Yêu cầu | Chọn |
|---|---|
| Windowed aggregation, sessionization, stream join, exactly-once | Managed Flink |
| ETL micro-batch vào data lake (Parquet/Iceberg), đội Spark | Glue streaming |
| Transform/filter/alert đơn giản, stateless | Lambda |
| Chỉ deliver + convert format | Firehose (kèm Lambda transform) |
| SQL tương tác khám phá stream | Managed Flink Studio notebook |
4. Kiến trúc mẫu: realtime dashboard
5. Điểm hay bị nhầm
- "Managed Flink" chạy liên tục và tính KPU-hour — job xử lý 1 lần/ngày thì đắt vô ích → dùng batch (Glue) cho việc đó
- Flink đọc từ KDS/MSK — nó không thay thế tầng ingest, nó là tầng processing
- Kinesis Data Analytics for SQL (sản phẩm cũ) đã ngừng nhận khách mới — đề hiện tại dùng Managed Flink (Flink SQL/Table API thay thế)
- Exactly-once của Flink là processing semantics nhờ checkpoint + sink hỗ trợ transactional/idempotent; sink không hỗ trợ thì vẫn có thể duplicate ở đầu ra
Câu hỏi ôn tập
-
Tính "số đơn hàng trung bình trong 10 phút gần nhất, cập nhật mỗi phút" — loại window nào?
Xem đáp án
Sliding (hopping) window: kích thước 10 phút, trượt mỗi 1 phút — các cửa sổ chồng lấp nhau, mỗi phút phát một kết quả cho 10 phút gần nhất. Tumbling chỉ phát mỗi 10 phút một lần và không chồng lấp — không thỏa "cập nhật mỗi phút".
-
Gom các event của người dùng thành "phiên", phiên kết thúc khi người dùng không hoạt động 30 phút. Window nào?
Xem đáp án
Session window với inactivity gap 30 phút — cửa sổ không có kích thước cố định, tự đóng khi không có event mới trong gap. Đây là dạng câu "sessionization" đặc trưng chỉ Flink/stream processor stateful làm gọn được.
-
Vì sao Lambda không phù hợp cho windowed aggregation 30 phút trên Kinesis?
Xem đáp án
Lambda là stateless per-invocation: window 30 phút cần giữ state giữa hàng nghìn invocation → phải tự lưu DynamoDB/ElastiCache, tự xử lý late data, khôi phục khi lỗi — phức tạp và dễ sai. (Lambda tumbling window có hỗ trợ tới 15 phút nhưng vẫn giới hạn.) Managed Flink sinh ra cho việc này: state managed, checkpoint, watermark, exactly-once.
-
Flink app bị lỗi và restart — dữ liệu đang xử lý có bị tính hai lần không?
Xem đáp án
Không (với cấu hình đúng): Flink checkpoint định kỳ state + vị trí đọc nguồn; khi restart, app khôi phục từ checkpoint gần nhất và đọc lại từ offset đó → exactly-once processing cho state nội bộ. Đầu ra cần sink transactional hoặc idempotent để end-to-end không duplicate.
-
Đội muốn khám phá dữ liệu stream bằng SQL tương tác trước khi viết app chính thức. Công cụ nào?
Xem đáp án
Managed Service for Apache Flink Studio — notebook Zeppelin managed, viết Flink SQL tương tác trực tiếp trên Kinesis/MSK, sau đó deploy notebook thành ứng dụng chạy liên tục.
-
Stream JSON từ Kinesis cần ghi thành bảng Iceberg trên S3 mỗi ~1 phút, đội quen PySpark. Lựa chọn hợp lý?
Xem đáp án
Glue streaming job (Spark Structured Streaming): đọc Kinesis, micro-batch ~1 phút, ghi Iceberg native (Glue 4.0/5.0), tận dụng kỹ năng PySpark sẵn có. Flink cũng làm được (sink Iceberg) nhưng đội không quen Java/Flink thì Glue streaming là "path of least resistance"; Firehose hiện ghi được Iceberg tables nhưng không chứa logic transform phức tạp.
Bài tập thực hành
- Viết pseudo-SQL Flink cho tumbling window 5 phút đếm event theo
event_type - Xác định loại window cho 5 yêu cầu tự nghĩ (mỗi loại ít nhất 1)
- So sánh chi phí: Flink 2 KPU chạy 24/7 vs Lambda 1M invocations/ngày (ước lượng thô)
- Đọc Flink windowing concepts
Tài liệu tham khảo chính thức
- Amazon Managed Service for Apache Flink
- Managed Flink Studio
- Glue streaming ETL jobs
- Lambda with Kinesis — windows
Tiếp theo: Quiz Tuần 3