Xử lý dữ liệu phân tán quy mô lớn với MapReduce
1. Tổng quan
A. Định nghĩa
MapReduce là mô hình lập trình và framework thực thi xử lý phân tán: chia dữ liệu đầu vào lớn thành các bản ghi khóa–giá trị, chạy
Mapsong song, gom các giá trị có cùng khóa trung gian và tổng hợp bằngReduce.
Nhà phát triển biểu diễn phép tính nghiệp vụ chủ yếu bằng hàm map và reduce, còn runtime xử lý việc chia input, lập lịch, truyền mạng, sắp xếp và khôi phục lỗi. Nhờ đó logic nghiệp vụ được tách khỏi độ phức tạp của cluster.
Mô hình được hệ thống hóa trong bài báo của Google năm 2004 về đơn giản hóa xử lý dữ liệu lớn và phổ biến qua Hadoop MapReduce mã nguồn mở. Spark, SQL phân tán và engine stream thường phù hợp hơn cho tính lặp hoặc độ trễ thấp, nhưng chia vùng, shuffle và tổng hợp vẫn là nguyên lý nền tảng.
B. Bối cảnh và sự cần thiết
Batch trên một máy chủ phải dựa vào mở rộng dọc khi dữ liệu vượt bộ nhớ, đĩa hoặc năng lực xử lý. Log, clickstream, cảm biến và giao dịch tăng nhanh khiến một máy lớn hơn không còn là lời giải kinh tế và an toàn.
Xử lý phân tán chia dữ liệu lên nhiều máy và chạy song song. Nếu tự viết chương trình phân tán, ứng dụng phải lo phân vùng, lỗi nút, retry, ghép kết quả và điều phối mạng. MapReduce đưa các mối quan tâm lặp lại đó vào framework.
Data locality rất quan trọng. Đưa mọi bản ghi về máy trung tâm tạo nút thắt mạng, còn MapReduce cố gắng đặt mapper gần nút đang giữ block đầu vào.
C. Mục tiêu chính
Mục tiêu không chỉ là dùng thêm máy chủ, mà là kết hợp khả năng mở rộng throughput, chịu lỗi, lập trình đơn giản, locality dữ liệu và tự động hóa vận hành.
Input độc lập có thể chạy song song, các bản ghi trung gian được phân phối lại theo khóa và chỉ task lỗi được chạy lại. Đổi lại, các giai đoạn trung gian thường được ghi xuống đĩa, nên workload tương tác hoặc lặp nhiều lần có thể cần engine khác.
2. Mô hình lập trình và cấu trúc thực thi
A. Trừu tượng khóa–giá trị
MapReduce xem input và output là các cặp <key, value>. Khóa có thể là offset file hoặc khóa DB; giá trị có thể là dòng văn bản, tài liệu JSON hoặc bản ghi cảm biến.
Hàm map đọc một bản ghi và phát ra không hoặc nhiều cặp trung gian. Hàm reduce nhận một khóa trung gian và danh sách giá trị gắn với khóa đó.
map(k1, v1) -> list(k2, v2)
reduce(k2, list(v2)) -> list(k3, v3)
Map lo biến đổi và lọc theo bản ghi, còn reduce lo tổng hợp và hợp nhất theo khóa logic. Ranh giới này làm rõ phần việc độc lập và phần việc theo nhóm.
B. Kiến trúc tổng thể
flowchart LR
IN[Tệp hoặc bảng đầu vào] --> SPLIT[Input Split]
SPLIT --> MR[Mapper]
MR --> COMB[Combiner tùy chọn]
COMB --> PART[Partitioner]
PART --> SHUF[Shuffle và Sort]
SHUF --> RED[Reducer]
RED --> OUT[Đầu ra hệ thống tệp phân tán]
RM[ResourceManager] -.lập lịch và tài nguyên.-> MR
RM -.lập lịch và tài nguyên.-> RED
NM[NodeManager] -.thực thi và báo trạng thái.-> MR
NM -.thực thi và báo trạng thái.-> RED
Trong Hadoop, ResourceManager điều phối tài nguyên và vị trí, còn NodeManager chạy container và báo trạng thái nút. ApplicationMaster theo dõi map/reduce task và yêu cầu retry.
InputFormat và RecordReader biến tệp thành bản ghi logic. Với văn bản, offset có thể là khóa và dòng là giá trị; bảo toàn ranh giới bản ghi là điều kiện đúng đắn.
Mapper đệm output trung gian trên lưu trữ cục bộ. Kết quả được chia partition, sắp xếp theo khóa và lưu cùng chỉ mục để reducer lấy phần của mình; đây là output tạm cho đến khi job commit.
C. Vòng đời xử lý
- Client gửi đường dẫn input/output, mapper, reducer, số partition và cấu hình.
- Framework chia input thành các split.
- Đặt map task, ưu tiên nút đang giữ dữ liệu.
- Mapper đọc bản ghi và phát ra cặp khóa–giá trị trung gian.
- Combiner tùy chọn thực hiện tổng hợp cục bộ.
- Partitioner gán mỗi khóa cho reducer.
- Reducer lấy các partition từ mapper và merge-sort.
- Reducer nhóm giá trị theo khóa và tính toán.
- Output format ghi kết quả và framework commit job.
Shuffle không chỉ là sao chép mạng. Nó lấy phần của reducer từ mọi mapper rồi merge theo khóa, nên thường chi phối thời gian và I/O.
sequenceDiagram
participant C as Client
participant AM as ApplicationMaster
participant M as Nút Mapper
participant R as Nút Reducer
C->>AM: gửi job (input, output, hàm)
AM->>M: đặt map task theo split
M->>M: tạo và sort cục bộ cặp trung gian
M->>R: shuffle dữ liệu partition
AM->>R: chạy reduce task
R->>R: nhóm và tổng hợp theo khóa
R-->>C: commit output và báo trạng thái
ApplicationMaster không xử lý bản ghi tập trung. Nó là control plane theo dõi vị trí, trạng thái và yêu cầu retry; mapper và reducer mới thực hiện phép tính dữ liệu.
3. Nguyên lý các thành phần Map và Reduce
A. Mapper
Mapper thường xử lý mỗi bản ghi độc lập. Ví dụ log có status 500 sẽ phát ra <service, 1>, còn bản ghi lỗi bị loại sớm để giảm dữ liệu shuffle.
Mapper cũng chịu trách nhiệm đổi kiểu, làm sạch dữ liệu và thiết kế khóa partition. Khóa kém có thể chia một nhóm nghiệp vụ thành nhiều khóa hoặc làm một reducer quá tải bởi khóa phổ biến.
B. Combiner
Combiner là phép tổng hợp cục bộ tùy chọn giữa mapper và reducer. Trong WordCount, 10.000 lần xuất hiện cloud tại một mapper có thể rút thành <cloud, 10000>.
Phép tính phải an toàn khi tổng hợp từng phần. Sum, count, min và max thường phù hợp; average cần truyền cả tổng và số lượng thay vì lấy trung bình của các trung bình.
Framework không bảo đảm combiner có chạy hay chạy bao nhiêu lần. Vì vậy không được đặt tính đúng đắn chỉ trong combiner.
C. Partitioner
Partitioner quyết định reducer sở hữu một khóa trung gian. Hash là cách phổ biến; ngày, vùng hoặc nhóm khách hàng có thể cần partitioner theo miền nghiệp vụ.
Mọi giá trị cùng khóa logic phải vào cùng reducer, đồng thời tải phải cân bằng. Sản phẩm, vùng hoặc ngày phổ biến có thể gây data skew.
Có thể thêm salt cho hot key, chia lên nhiều partition rồi tổng hợp lần hai. Cách này thêm một bước merge nên phải đánh giá cả tính đúng và chi phí.
D. Shuffle và Sort
Shuffle chuyển partition từ mapper tới reducer, tiêu thụ băng thông và I/O đĩa, thường là nguồn trễ lớn nhất.
Sort và merge làm cho các khóa giống nhau đứng liên tiếp. Reducer có thể stream từng nhóm khóa thay vì nạp cả dataset vào bộ nhớ.
Lọc sớm, combiner an toàn, nén intermediate và chọn số partition phù hợp là tối ưu cơ bản. Tuy nhiên nén có thể chậm hơn nếu CPU đắt hơn phần mạng tiết kiệm được.
E. Reducer
Reducer tổng hợp, sắp xếp, join, loại trùng hoặc tính thống kê cho một khóa. Vì input đã theo thứ tự khóa, state có thể giải phóng khi khóa thay đổi.
Reducer có thể bị chạy lại. Ghi ra hệ thống bên ngoài cần idempotent key, output tạm hoặc commit nguyên tử để retry không tạo bản ghi hay lời gọi trùng.
4. Chịu lỗi và thiết kế hiệu năng
A. Chạy lại task
Ở quy mô cluster, hỏng đĩa, đứt mạng, dừng process và node quá tải là bình thường. MapReduce chạy lại task lỗi trên node khác thay vì chạy lại toàn bộ batch.
Block input được nhân bản giúp mapper dùng bản sao khác. Với straggler, speculative execution có thể chạy cùng task ở node thứ hai và nhận kết quả thành công trước.
Side effect không idempotent như thanh toán bên ngoài rất nguy hiểm khi chạy trùng. Mapper và reducer nên là hàm thuần hoặc dùng giao thức output idempotent.
B. Locality và bố trí tệp
Hiệu năng thường bị giới hạn bởi di chuyển dữ liệu hơn là CPU. Đọc block cục bộ tránh mạng, còn đặt task ở node xa làm tăng chi phí truyền.
Quá nhiều tệp nhỏ tạo overhead split và khởi động task; quá ít split lớn làm giảm song song. Kích thước tệp, block, số task và số node phải được điều chỉnh cùng nhau.
C. Chi phí và quan sát
Gọi chi phí tính map là (M), số bản ghi trung gian là (I), số reducer là (R), tổng chi phí gồm map, ghi intermediate, shuffle/sort của (I) và reduce.
Theo dõi input/output của map, tỷ lệ giảm bởi combiner, byte shuffle, thời gian chờ, chênh lệch input giữa reducer, số lỗi và retry. Một reducer chậm gợi ý skew; mọi reducer cùng chờ gợi ý áp lực mạng hoặc đĩa.
Với job lặp lại, quản lý p95 và p99 cùng baseline về volume, không chỉ nhìn thời gian trung bình.
5. Ví dụ và ứng dụng
A. WordCount
Mapper phát <word, 1>, combiner cộng cục bộ và reducer cộng mọi phần để tạo <word, total>. Bài học chính là các khóa giống nhau từ nhiều node vẫn được gom thành một nhóm logic.
B. Inverted Index
Mapper phát <word, documentId>, reducer sắp xếp và loại trùng ID để tạo posting list. Từ phổ biến có thể tập trung lượng giá trị lớn, nên cần loại stop word, chia danh sách nóng hoặc merge bổ sung.
C. Tổng hợp log và giao dịch
Đếm lỗi theo ngày và dịch vụ dùng <date|service, 1>, tổng tiền theo khách hàng dùng <customerId, amount>. Dữ liệu bị quản lý vẫn cần mã hóa, masking, quyền truy cập, audit và chính sách lưu giữ.
6. So sánh với công nghệ thay thế
MapReduce mạnh ở batch dựa trên đĩa và retry tự động, nhưng ghi intermediate nhiều lần làm workload lặp chậm. Việc chọn phải dựa vào latency, state, tái sử dụng, dạng dữ liệu và SLA.
| Tiêu chí | MapReduce | Spark | SQL phân tán | Engine stream |
|---|---|---|---|---|
| Cơ bản | Batch theo giai đoạn | DAG cache và batch | Truy vấn khai báo | Event liên tục |
| Trung gian | Thường ở đĩa | Cache hoặc memory | Do engine lập kế hoạch | State và checkpoint |
| Ưu điểm | Đơn giản, retry | Lặp và latency thấp | Năng suất SQL | Cửa sổ thời gian thực |
| Hạn chế | Shuffle và đĩa | Tuning bộ nhớ | Giới hạn UDF | Quản lý state, trùng |
Spark mở rộng ý tưởng map/reduce bằng RDD, DataFrame và DAG, cho phép cache dữ liệu dùng lại; áp lực bộ nhớ và chi phí shuffle vẫn tồn tại.
SQL phân tán có thể tự tối ưu join, filter pushdown và partition pruning. MapReduce vẫn tự nhiên hơn với logic tùy biến hoặc bước tích hợp hệ thống ngoài.
Stream engine xử lý input vô hạn bằng window và state, nên phải định nghĩa event time, watermark, late event, duplicate và guarantee riêng.
7. Đào sâu — vị trí trong nền tảng hiện đại và hướng thi
Nên hiểu MapReduce như cách tư duy về dataflow phân tán hơn là chỉ một sản phẩm. Chia input, phân phối lại theo khóa, tổng hợp phần, sort và merge cuối đều xuất hiện trong kế hoạch SQL phân tán.
Khi object storage và container tách rời lưu trữ với tính toán, locality chuyển từ đĩa cục bộ sang định dạng tệp, partition pruning, cache và tối ưu chi phí mạng.
Trong bài thi, nên vẽ Map → Combine → Partition → Shuffle/Sort → Reduce, dùng WordCount hoặc tổng hợp log để giải thích gom nhóm theo khóa, rồi nối với locality, retry, skew và bottleneck shuffle.
8. Lưu ý và hàm ý
- Kiểm tra bài toán có hợp với gom nhóm khóa–giá trị không. Join nhiều tầng và intermediate phình to có thể cần SQL, graph hoặc stream.
- Quản lý shuffle như trung tâm chi phí. Đo volume trung gian, độ lệch reducer, mạng và đĩa trước khi tăng cluster.
- Coi combiner là tùy chọn. Chỉ đặt phép tổng hợp từng phần an toàn và giữ tính đúng ở reducer.
- Làm retry an toàn. Dùng khóa idempotent, đường dẫn tạm và commit nguyên tử cho side effect.
- Xử lý skew ở cấp khóa nghiệp vụ. Tăng số reducer không tự giải quyết một hot key.
- Thiết kế tệp và split cùng nhau. Gom tệp nhỏ, giữ song song phù hợp và xác minh partition pruning giảm lượng đọc.
- Đưa governance vào đường chạy. Mã hóa dữ liệu nhân bản và trung gian, kiểm soát quyền, masking dữ liệu cá nhân và ghi audit.
- Chọn theo SLA và chi phí. So sánh thời gian, thời gian phục hồi, chi phí cluster, độ khó vận hành và năng lực đội ngũ với Spark, SQL, stream.
Tài liệu tham khảo
- Dean, J. and Ghemawat, S., “MapReduce: Simplified Data Processing on Large Clusters”, Google Research. https://research.google/pubs/mapreduce-simplified-data-processing-on-large-clusters/
- Apache Hadoop, “MapReduce Tutorial”. https://hadoop.apache.org/docs/current/hadoop-mapreduce-client/hadoop-mapreduce-client-core/MapReduceTutorial.html
- Apache Hadoop, “HDFS Architecture Guide”. https://hadoop.apache.org/docs/current/hadoop-project-dist/hadoop-hdfs/HdfsDesign.html
Tóm tắt một câu: MapReduce mở rộng batch lớn bằng cách chia input khóa–giá trị và chạy Map, Combine, Shuffle/Sort, Reduce song song, kết hợp locality dữ liệu, retry task và tổng hợp theo khóa để đạt khả năng mở rộng và chịu lỗi.