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

Tuần 3 - Ngày 5: Managed Service for Apache Flink và Streaming ETL

Tuần 3 – Ngày 5

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

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.

SOURCESFLINKAPPSINKSKinesisDataStreamsKinesisDataStreamsMSK/Kafkaoperators:FirehoseS3filter,map,S3,DynamoDB,join,WINDOW,OpenSearch,...aggregate(sinkconnectors)[STATE+CHECKPOINT]Viếtbng:Java/Scala/Python(TableAPI/DataStreamAPI/SQL)Studionotebook(Zeppelin)chopháttrintươngtác

Đặ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ĩaVí dụ
TumblingCửa sổ cố định, không chồng lấpDoanh 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á gapPhiê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ầuChọn
Windowed aggregation, sessionization, stream join, exactly-onceManaged Flink
ETL micro-batch vào data lake (Parquet/Iceberg), đội SparkGlue streaming
Transform/filter/alert đơn giản, statelessLambda
Chỉ deliver + convert formatFirehose (kèm Lambda transform)
SQL tương tác khám phá streamManaged Flink Studio notebook

4. Kiến trúc mẫu: realtime dashboard

PaymenteventsKDSManagedFlinktumblingwindow1phút:tnggiaodch,đếmlipháthinbtthưng(songưng/State)CloudWatchcustommetrics/DynamoDB/KDSkhácDashboard+CloudWatchAlarmSNS(songsong)KDSFirehoseS3(ngunstht,phântíchsau)

5. Điểm hay bị nhầm

  1. "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 đó
  2. Flink đọc từ KDS/MSK — nó không thay thế tầng ingest, nó là tầng processing
  3. 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ế)
  4. 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

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

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

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

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

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

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


Tiếp theo: Quiz Tuần 3