Beyond MapReduce — Spark, Flink và thế hệ data processing mới (DDIA)

Phong

Mở đầu

MapReduce đã từng là một cuộc cách mạng trong thế giới xử lý dữ liệu lớn. Google công bố paper năm 2004, và nhanh chóng trở thành tiêu chuẩn cho batch processing tại quy mô internet. Nhưng công nghệ không ngừng phát triển, và MapReduce — dù mạnh mẽ — có những hạn chế cố hữu mà các thế hệ sau đã giải quyết triệt để.

Chương 10 của DDIA dành phần lớn để phân tích MapReduce, nhưng phần cuối (10.2) bàn về những giới hạn đó và các hệ thống thay thế hiện đại hơn như Spark, Flink, và Tez.

Ảnh: Wolfgang Weiser — Pexels

Vấn Đề Lớn Nhất Của MapReduce: Materialized Intermediates

Điểm yếu chí tử của MapReduce là nó ghi kết quả trung gian (intermediate state) xuống disk giữa mỗi phase. Một job MapReduce điển hình có thể có nhiều vòng map→reduce, và mỗi vòng đều phải:

  • Ghi output của reduce vào distributed filesystem (HDFS)
  • Job tiếp theo đọc lại từ HDFS
  • Lặp lại — mỗi bước là một lần I/O đắt đỏ

Hãy tưởng tượng bạn làm bánh: mỗi công đoạn bạn lại bỏ bánh vào tủ lạnh, hôm sau lấy ra làm tiếp, rồi lại bỏ vào tủ lạnh. MapReduce làm điều đó với dữ liệu — và trên cluster hàng trăm node, chi phí I/O này cực kỳ tốn kém.

Ngoài ra, mỗi MapReduce job chạy riêng rẽ — không có cơ chế pipeline giữa các tác vụ. Nếu bạn muốn thực hiện 3 bước xử lý nối tiếp, bạn phải viết 3 job MapReduce, mỗi job tốn vài phút để khởi tạo và ghi kết quả trung gian.

Dataflow Engines — Spark, Flink, Tez

Các hệ thống thế hệ tiếp theo (gọi chung là dataflow engines) giải quyết vấn đề này bằng một ý tưởng đơn giản: thay vì ghi mọi kết quả trung gian xuống disk, hãy pipe dữ liệu trực tiếp từ operator này sang operator khác trong cùng một process.

Apache Spark là đại diện nổi bật nhất. Spark dùng khái niệm Resilient Distributed Datasets (RDDs) — một collection phân tán, bất biến, có thể được tái tính toán nếu mất dữ liệu. Thay vì ghi intermediate state ra HDFS, Spark giữ dữ liệu trong memory (hoặc spill ra disk nếu cần) và xây dựng một DAG (Directed Acyclic Graph) các transformations.

Apache Flink đi xa hơn: nó không chỉ xử lý batch mà còn xử lý streaming với event-by-event processing. Flink coi batch như một trường hợp đặc biệt của streaming (bounded stream) — cùng một engine xử lý cả hai workload.

Ảnh: Jakub Zerdzicki — Pexels

Apache Tez (dùng trong Hive, Pig) cho phép xây dựng DAG phức tạp từ nhiều bước xử lý, tối ưu hoá việc tái sử dụng dữ liệu trung gian.

Cả ba đều có điểm chung: thay thế mô hình map→reduce→disk→map→reduce→disk bằng một DAG operators linh hoạt, nơi dữ liệu được truyền trực tiếp giữa các bước.

Fault Tolerance — Recomputation vs Checkpointing

Một trong những ưu điểm của MapReduce là fault tolerance đơn giản: nếu một task thất bại, nó chỉ cần chạy lại task đó trên cùng một input (vì input đã được ghi trong HDFS và intermediate state không được chia sẻ).

Dataflow engines có hai cách tiếp cận:

  • Recomputation (Spark): Dựa vào lineage — mỗi RDD biết nó được tạo ra từ RDD nào. Nếu một partition mất, Spark chỉ cần tính lại partition đó từ RDD gốc. Cách này hiệu quả nếu DAG không quá sâu.
  • Checkpointing (Flink): Định kỳ ghi snapshot trạng thái của toàn bộ streaming pipeline vào storage bền vững. Khi có lỗi, rollback về checkpoint gần nhất. Cách này phù hợp hơn với streaming workloads, nơi không có "input gốc" để recompute.

Điểm thú vị: Spark ban đầu chỉ dùng recomputation, nhưng với các DAG sâu hoặc các transformation đắt đỏ, nó cũng cho phép checkpointing như một tuỳ chọn.

Sorting — MapReduce's Obsession

Một điểm đặc biệt của MapReduce là nó sort dữ liệu giữa map và reduce (shuffle phase). Việc sort này có ích trong nhiều trường hợp, nhưng cũng là một chi phí lớn.

Dataflow engines linh hoạt hơn: họ chỉ sort khi cần thiết (join, groupBy) và có thể dùng các cấu trúc dữ liệu khác như hash tables cho những tác vụ không yêu cầu thứ tự.

Ảnh: ThisIsEngineering — Pexels

Từ Batch Sang Streaming — Sự Hội Tụ

Một xu hướng quan trọng trong chapter này là sự hội tụ giữa batch và streaming processing.

MapReduce chỉ xử lý batch — bạn nạp toàn bộ dữ liệu, xử lý, ra kết quả. Nhưng thực tế, dữ liệu thường đến liên tục (logs, events, sensor data). Các hệ thống thế hệ mới như Spark Streaming, Flink, Kafka Streams cho phép xử lý dữ liệu với độ trễ thấp hơn nhiều — từ vài giây đến milliseconds.

DDIA gọi đây là "the unification of batch and stream processing" — một engine có thể xử lý cả batch lẫn streaming, giúp đơn giản hoá kiến trúc hệ thống.

Key Takeaways

  • Materialized intermediates là điểm yếu lớn nhất của MapReduce — ghi mọi thứ ra disk giữa các phase
  • Dataflow engines (Spark, Flink, Tez) pipe dữ liệu trực tiếp giữa các operators, tránh I/O không cần thiết
  • Spark dùng RDDs + DAG + recomputation để đạt fault tolerance
  • Flink xử lý streaming true event-by-event, checkpointing để phục hồi
  • Xu hướng hội tụ: một engine xử lý cả batch lẫn streaming
  • Không phải lúc nào cũng cần sort — dataflow engines dùng hash tables khi có thể

📋 Phụ lục thuật ngữ

Thuật ngữÝ nghĩa
Materialized IntermediateKết quả trung gian được ghi xuống disk giữa các phase xử lý
Dataflow EngineHệ thống xử lý dữ liệu dạng DAG, pipe dữ liệu giữa các operators
RDDResilient Distributed Dataset — collection phân tán, bất biến, có lineage để recompute
DAGDirected Acyclic Graph — biểu diễn luồng xử lý dữ liệu
LineageLịch sử biến đổi của dữ liệu — cho phép tái tính toán khi mất partition
CheckpointingSnapshot định kỳ trạng thái hệ thống để rollback khi có lỗi
ShuffleQuá trình phân phối lại dữ liệu giữa các partition (thường kèm sort)

Kết

MapReduce đã mở đường, nhưng thế hệ dataflow engines như Spark và Flink mới thực sự đưa xử lý dữ liệu lớn lên một tầm cao mới. Bằng cách loại bỏ materialized intermediates, pipeline dữ liệu giữa các bước, và hỗ trợ cả batch lẫn streaming, các hệ thống này giải quyết được những hạn chế cốt lõi của MapReduce.

Bài tiếp theo trong series sẽ bước sang Chương 11 — Stream Processing, nơi đi sâu vào event streams, CDC, và event sourcing.