Xử lý dữ liệu theo lô — Unix Philosophy & MapReduce (Designing Data-Intensive Applications)

Phong

Mở đầu

Chương 10 của Designing Data-Intensive Applications bàn về một chủ đề tưởng cũ mà không hề cũ: batch processing (xử lý dữ liệu theo lô). Trong thời đại real-time streaming lên ngôi, nhiều dev có thể nghĩ batch processing là thứ "xưa rồi". Nhưng thực tế, batch vẫn là xương sống của biết bao hệ thống lớn — từ index search engines, build recommendation systems, đến ETL pipelines cho data warehouse.

Điều thú vị là tác giả bắt đầu chương này từ một thứ rất đơn giản: Unix pipes. Và từ đó, ông chỉ ra rằng MapReduce thực chất là... Unix pipes cho distributed systems. Nghe có vẻ đơn giản, nhưng để hiểu được sự tương đồng đó, bạn cần nắm được triết lý đằng sau.

Ảnh: Pixabay — Pexels

Triết lý Unix — nhỏ mà có võ

Doug McIlroy, cha đẻ của Unix pipes, đã đúc kết triết lý Unix như sau:

  • Làm một việc và làm thật tốt
  • Kết nối các chương trình nhỏ qua pipes để giải quyết bài toán phức tạp
  • Dùng text làm universal interface — bất kỳ chương trình nào cũng đọc được

Nghe quen không? Đây chính là tư tưởng modular design — mỗi component làm một việc, giao tiếp qua interface chuẩn. Trong thế giới Unix, interface đó là stdin/stdout với text stream.

Xử lý log với Unix pipes

Giả sử bạn có một file log web server và muốn tìm top 5 URL được truy cập nhiều nhất:

cat access.log | awk '{print $7}' | sort | uniq -c | sort -rn | head -5

Một dòng lệnh duy nhất, không cần code phức tạp. Mỗi tool làm một bước:

  • cat — đọc file
  • awk — trích cột URL (field 7)
  • sort — sắp xếp URL theo thứ tự alphabet
  • uniq -c — gộp URL trùng và đếm
  • sort -rn — sắp xếp theo số lượt giảm dần
  • head -5 — lấy top 5

Đây là batch processing thuần tuý: dữ liệu được đọc, xử lý qua từng bước, và cho ra kết quả. Mỗi bước chạy tuần tự, nhưng nhờ pipe, bước trước bắt đầu output ngay khi có dữ liệu — không cần đợi xong hết mới chuyển sang bước sau.

Ảnh: Tima Miroshnichenko — Pexels

Từ Unix pipes đến MapReduce

Unix pipes rất mạnh, nhưng có một giới hạn: tất cả chạy trên một máy. Khi bạn có terabytes dữ liệu, Unix pipes không đủ — bạn cần distributed processing. Và đó là lúc MapReduce xuất hiện.

Jeffrey Dean và Sanjay Ghemawat (Google) đã thiết kế MapReduce như một Unix pipe cho distributed systems. Cùng một tư tưởng: chain các xử lý nhỏ lại với nhau, nhưng chạy trên hàng nghìn máy song song.

MapReduce gồm hai phase chính:

Map phase

Mỗi worker đọc một phần dữ liệu (thường là một block HDFS, thường 64-128MB), áp dụng hàm map do người dùng định nghĩa, và emit ra các cặp key-value. Bước này song song hoá hoàn toàn — các worker không phụ thuộc nhau.

Ví dụ kinh điển: đếm số lần xuất hiện của mỗi từ trong một bộ sưu tập document. Hàm map nhận một document, split thành các từ, emit (word, 1) cho mỗi từ.

Reduce phase

Sau khi map xong, hệ thống shuffle tất cả cặp key-value — gom các cặp có cùng key về cùng một worker. Worker reduce nhận một key và danh sách tất cả values của key đó, áp dụng hàm reduce, và output kết quả cuối cùng.

Với bài toán đếm từ, reduce nhận ("hello", [1, 1, 1, ...]) và output ("hello", tổng_số).

Sự khác biệt giữa Unix pipes và MapReduce

Dù cùng tư tưởng, MapReduce có vài điểm khác biệt quan trọng:

  • Thứ tự xử lý: Unix pipes xử lý tuần tự từng bước, mỗi bước chạy hết dữ liệu rồi mới chuyển (dù pipe cho phép streaming). MapReduce có hai phase riêng biệt — map hoàn toàn xong mới bắt đầu shuffle và reduce.
  • Sort là trung tâm: Trong MapReduce, sort là thao tác cốt lõi. Sau map, các cặp key-value được sort theo key trong quá trình shuffle. Trong Unix pipes, sort là optional — bạn chỉ dùng khi cần.
  • Input/output: Unix pipes đọc từ stdin, ghi ra stdout — dữ liệu là text stream. MapReduce đọc/ghi từ distributed filesystem (GFS/HDFS) — tiêu chuẩn là các file cố định.
  • Fault tolerance: Unix pipe hỏng một lệnh là cả pipe hỏng. MapReduce tự động retry task thất bại — một trong những đóng góp lớn nhất của MapReduce là khả năng chạy ổn định trên hàng nghìn máy mà không sợ failure.
Ảnh: Markus Winkler — Pexels

MapReduce workflows phức tạp hơn

Không phải bài toán nào cũng chỉ cần một MapReduce. Nhiều bài toán cần chuỗi MapReduce nối tiếp nhau — output của job này là input của job kế tiếp. Điều này rất giống với Unix pipes: bạn chain nhiều câu lệnh với pipes.

Ví dụ trong machine learning: tính ma trận tương đồng giữa các user (cho collaborative filtering) có thể cần 4-5 MapReduce jobs: tính co-occurrence → tính tần suất → chuẩn hoá → ranking. Mỗi job viết output ra HDFS, job sau đọc vào.

Các framework như Pig (Yahoo!) và Hive (Facebook) ra đời để abstract hoá những workflow này — bạn viết script ngắn (Pig Latin, HiveQL), framework tự sinh ra chuỗi MapReduce jobs phía sau.

Key Takeaways

  • Unix Philosophy — modular design với pipes và text interface là nền tảng của batch processing
  • MapReduce = Unix pipes for distributed systems — cùng tư tưởng chain các xử lý nhỏ, nhưng chạy song song trên cluster
  • Map phase — biến đổi dữ liệu đầu vào thành cặp key-value, song song hoá hoàn toàn
  • Reduce phase — gom nhóm theo key và tổng hợp, cần shuffle/sort giữa hai phase
  • Fault tolerance — retry task tự động giúp MapReduce chạy ổn định trên hàng nghìn máy
  • Sort là cốt lõi — khác với Unix pipes (sort optional), MapReduce built-in sort cho shuffle phase

Kết

Một trong những điểm mạnh của DDIA là cách tác giả kết nối những ý tưởng tưởng chừng cũ kỹ (Unix pipes từ những năm 70) với các hệ thống hiện đại (MapReduce, Spark). Triết lý Unix vẫn còn nguyên giá trị: small tools, composable, text interface — và MapReduce là một minh chứng cho việc áp dụng triết lý đó ở quy mô hoàn toàn khác.

Ở chương tiếp theo, tác giả sẽ bàn về stream processing — một cách tiếp cận khác cho xử lý dữ liệu, nơi delay được giảm xuống mức tối thiểu. Nếu bạn muốn hiểu sự khác nhau giữa batch và stream, chương 11 là không thể bỏ qua.

📋 Phụ lục thuật ngữ

  • Batch processing — xử lý dữ liệu theo lô, input cố định, output sinh ra sau khi xử lý xong
  • Unix pipe — cơ chế kết nối stdout của process này với stdin của process khác
  • MapReduce — programming model của Google cho xử lý dữ liệu song song trên cluster
  • Shuffle — giai đoạn trung gian giữa map và reduce, nơi các cặp key-value được sort và gom nhóm
  • HDFS — Hadoop Distributed File System, lưu trữ dữ liệu phân tán cho MapReduce
  • Fault tolerance — khả năng hệ thống tiếp tục hoạt động khi có component bị lỗi