Chương 10 — Batch Processing
Phần III — Derived Data

Chương 10 — Batch Processing

15 phút đọc DDIA · Martin Kleppmann

🎯 Mục tiêu chương: Hiểu batch processing — xử lý một tập dữ liệu đầu vào bounded (kích thước cố định, đã biết) để sinh ra output dẫn xuất — từ Unix pipeline tới MapReduce và các dataflow engine (Spark, Flink, Tez). Nắm các thuật toán join/grouping phân tán, cách chịu lỗi, và vì sao triết lý "input immutable, output thay thế hoàn toàn" giúp hệ dễ bảo trì.

Ba loại hệ thống

LoạiCách hoạt độngThước đo chính
Services (online)Chờ request, xử lý nhanh, trả response; thường có người dùng đang chờResponse time, availability
Batch processing (offline)Nhận input lớn, chạy job (vài phút tới vài ngày), sinh output; thường chạy theo lịchThroughput
Stream processing (near-real-time)Như batch nhưng xử lý event ngay sau khi chúng xảy ra, input không bao giờ kết thúcLatency thấp hơn batch (Ch11)

MapReduce (Google, 2004) từng được tung hô là "thuật toán làm Google scale được"; về mặt tư tưởng nó khá low-level so với các hệ song song của data warehouse trước đó, nhưng là bước tiến lớn về quy mô trên phần cứng phổ thông. Tầm quan trọng của MapReduce đang giảm, nhưng nó là công cụ học rất rõ ràng. Batch processing thực ra rất cũ: máy Hollerith trong điều tra dân số Mỹ năm 1890, máy phân loại thẻ đục lỗ IBM những năm 1940–50 giống MapReduce một cách kỳ lạ.

Batch Processing with Unix Tools

Simple Log Analysis

Bài toán: từ access log của nginx, tìm 5 trang được truy cập nhiều nhất.

cat /var/log/nginx/access.log |
  awk '{print $7}' |   # lấy trường thứ 7 = URL
  sort             |   # sắp xếp, các URL giống nhau đứng cạnh nhau
  uniq -c          |   # đếm các dòng liền kề giống nhau
  sort -r -n       |   # sắp theo số đếm, giảm dần
  head -n 5            # lấy 5 dòng đầu
  • Xử lý được hàng GB log trong vài giây; sửa phân tích rất dễ (bỏ file CSS: '$7 !~ /\.css$/ {print $7}'; đếm theo IP: '{print $1}').
  • awk, sed, grep, sort, uniq, xargs giải quyết được một lượng lớn bài toán phân tích.

Chain of commands vs custom program

Viết bằng Ruby: dùng một hash table counts[url] += 1, rồi sort theo count lấy top 5. Khác biệt không nằm ở cú pháp mà ở execution flow.

Sorting vs in-memory aggregation

CáchƯu điểmHạn chế
In-memory hash table (Ruby)Nhanh khi working set (số URL khác nhau) vừa memory; 1 triệu lượt cho 1 URL vẫn chỉ tốn 1 entryHỏng khi số key khác nhau vượt RAM
Sort (Unix)Dữ liệu lớn hơn memory vẫn chạy: sort từng chunk trong memory, ghi segment ra disk, rồi merge (giống SSTable/LSM) — truy cập tuần tự, thân thiện với diskTốn thêm công sort

GNU sort tự spill ra disk và song song hoá trên nhiều core → pipeline Unix scale tốt; nút cổ chai thường là tốc độ đọc file từ disk.

The Unix Philosophy

Doug McIlroy (người phát minh pipe, 1964): nối các chương trình như nối ống tưới vườn. Các nguyên tắc (1978): 1. Mỗi chương trình làm một việc cho tốt; việc mới thì viết mới thay vì nhồi feature. 2. Kỳ vọng output của mọi chương trình sẽ thành input của một chương trình khác, chưa biết trước. Không làm rối output, tránh format cột cứng nhắc hay binary, không đòi input tương tác. 3. Thiết kế để thử sớm (trong vài tuần), sẵn sàng vứt bỏ phần vụng về. 4. Dùng tool để giảm công việc, kể cả phải tự viết tool rồi vứt đi.

→ Nghe rất giống Agile/DevOps ngày nay. sort là ví dụ "làm một việc cho tốt" (tốt hơn sort trong thư viện chuẩn của đa số ngôn ngữ), nhưng chỉ mạnh khi kết hợp với uniq… Vậy điều gì cho phép composability?

A uniform interface

  • Mọi chương trình dùng cùng một interface: file (file descriptor) = chuỗi byte có thứ tự. File thật, pipe, socket, device driver (/dev/audio) đều dùng chung interface này.
  • Quy ước: ASCII text, record phân tách bằng \n. Việc parse field thì mơ hồ hơn (whitespace, tab, CSV…) — {print $7} không đẹp bằng {print $request_url}.
  • Tương tự: URL và HTTP là uniform interface của web (so với thời BBS phải quay số từng hệ riêng).
  • Ngày nay việc các phần mềm phối hợp mượt như Unix là ngoại lệ; ngay cả DB cùng data model cũng khó chuyển dữ liệu qua lại → "Balkanization" của dữ liệu.

Separation of logic and wiring

  • Chương trình chỉ đọc stdin, ghi stdout; người dùng shell quyết định nối vào đâu → loose coupling / late binding / inversion of control.
  • Tự viết tool (ví dụ chuyển IP thành mã quốc gia) rồi cắm vào pipeline được ngay.
  • Giới hạn: nhiều input/output thì khó; không pipe thẳng ra network được; nếu chương trình tự mở file/network thì mất tính linh hoạt của việc nối dây.

Transparency and experimentation

  • Input được coi là immutable → chạy lại bao nhiêu lần cũng được.
  • Có thể dừng pipeline ở bất kỳ đâu, pipe vào less để xem.
  • Ghi output một stage ra file để chạy lại stage sau mà không chạy lại cả pipeline.
  • Giới hạn lớn nhất: chỉ chạy trên một máy → đó là lý do cần Hadoop.

MapReduce and Distributed Filesystems

MapReduce giống Unix tools nhưng phân tán trên hàng nghìn máy: thô, brute-force nhưng hiệu quả. Một job ≈ một process Unix: không sửa input, không side effect ngoài việc ghi output; output ghi một lần, tuần tự.

HDFS (bản open source của Google File System): - Shared-nothing — khác shared-disk (NAS/SAN) cần phần cứng đặc biệt như Fibre Channel. - Mỗi máy chạy một daemon phục vụ file trên disk local; NameNode trung tâm theo dõi block nào nằm ở máy nào → một filesystem lớn gộp disk của mọi máy. - Chịu lỗi bằng replication hoặc erasure coding (Reed–Solomon, tốn ít dung lượng hơn) — giống RAID nhưng qua mạng datacenter thường. - Scale tới hàng chục nghìn máy, hàng trăm PB, rẻ hơn nhiều so với storage appliance chuyên dụng. - Tương tự: GlusterFS, QFS, object store (S3, Azure Blob, Swift). Khác biệt: HDFS cho phép chạy compute ngay trên máy chứa dữ liệu; object store thường tách storage và compute.

MapReduce Job Execution

Bốn bước, ánh xạ trực tiếp từ ví dụ log: 1. Đọc input, tách thành records (mỗi dòng log) — do input format parser làm. 2. Mapper: trích key/value từ mỗi record (awk '{print $7}' → key = URL). 3. Sort theo key (sort) — ngầm định, framework tự làm. 4. Reducer: duyệt các cặp key-value đã sort, gộp các giá trị cùng key (uniq -c).

  • Mapper: gọi một lần cho mỗi record, sinh 0..n cặp key-value, không giữ state giữa các record.
  • Reducer: nhận một key và iterator trên mọi value của key đó, sinh output.
  • Sort lần hai (xếp hạng URL) = cần một MapReduce job thứ hai. Vai trò: mapper chuẩn bị dữ liệu cho việc sort; reducer xử lý dữ liệu đã sort.

Distributed execution of MapReduce

  • Song song hoá dựa trên partitioning: mỗi file/block của input → một map task.
  • "Putting the computation near the data": scheduler cố chạy mapper trên máy có replica của block đó → giảm tải mạng. Framework copy code (ví dụ JAR) tới máy trước.
  • Số reduce task do người viết job cấu hình; framework dùng hash của key để chọn reducer → mọi cặp cùng key về cùng reducer.
  • Sort theo từng giai đoạn: mỗi mapper chia output theo reducer, ghi thành file đã sort trên disk local của mapper (kiểu SSTable).
  • Mapper xong → reducer kéo file thuộc partition của mình từ mọi mapper. Quá trình partition + sort + copy này gọi là shuffle (tên dễ nhầm — hoàn toàn không có ngẫu nhiên).
  • Reducer merge các file đã sort (giữ thứ tự) → các record cùng key đứng liền nhau; output ghi lên HDFS (một bản local + các replica).
  • Mapper/reducer trong Hadoop là class Java; trong MongoDB/CouchDB là function JavaScript; cũng có thể dùng Unix tools (Hadoop Streaming).

MapReduce workflows

  • Một job giải được ít bài toán → chain các job thành workflow: output directory của job trước = input directory của job sau (nối bằng tên thư mục; Hadoop không hỗ trợ workflow sẵn).
  • Giống ghi output của từng lệnh ra file tạm, không giống pipe.
  • Output chỉ hợp lệ khi job thành công hoàn toàn (output dở dang bị bỏ) → job sau chỉ chạy khi job trước xong. Có các workflow scheduler: Oozie, Azkaban, Luigi, Airflow, Pinball.
  • Hệ recommendation thường có 50–100 job. Các tool cấp cao hơn (Pig, Hive, Cascading, Crunch, FlumeJava) tự nối nhiều stage.

Reduce-Side Joins and Grouping

MapReduce không có index: đọc toàn bộ input = full table scan. Tệ nếu chỉ cần vài record, nhưng hợp lý cho analytics trên lượng lớn record và song song hoá được. "Join" trong batch = giải quyết mọi lần xuất hiện của một liên kết trong dataset (xử lý mọi user cùng lúc), không phải tra một user.

Example: analysis of user activity events

  • Bên trái: activity events/clickstream (fact table), chỉ chứa user ID. Bên phải: user database (dimension) chứa ngày sinh. Mục tiêu: trang nào phổ biến với nhóm tuổi nào.
  • Cách ngây thơ: với mỗi event, query DB từ xa → bị giới hạn bởi round-trip, cache phụ thuộc phân bố dữ liệu, dễ làm quá tải DB, và nondeterministic (DB thay đổi trong lúc chạy).
  • Cách đúng: lấy bản copy của user DB (qua ETL từ backup) đưa vào HDFS cùng với log, rồi để MapReduce gom dữ liệu liên quan vào một chỗ.

Sort-merge joins

  • Một nhóm mapper đọc events → emit (userID, event); nhóm khác đọc user DB → emit (userID, dateOfBirth).
  • Shuffle làm các record cùng userID đứng cạnh nhau ở reducer. Secondary sort: sắp để reducer luôn thấy record user trước, rồi tới các event theo timestamp.
  • Reducer: lưu ngày sinh vào biến local, duyệt các event, emit (viewed-url, viewer-age). Job sau tính phân bố tuổi theo URL.
  • Chỉ giữ một user record trong memory, không gọi network → nhanh.
  • Mapper như "gửi message" tới reducer; key đóng vai trò địa chỉ đích.
  • Tách phần network communication (đưa dữ liệu tới đúng máy) khỏi application logic; framework tự retry task lỗi → code ứng dụng không phải lo partial failure.

GROUP BY

  • Cũng là pattern "gom về cùng chỗ": mapper emit với key = grouping key. Aggregate: COUNT(*), SUM(field), top-k.
  • Sessionization: gom mọi activity event của một user session (rải rác trên log của nhiều web server) theo session cookie/user ID → phân tích hành vi, A/B testing, đánh giá marketing.

Handling skew

  • Hot keys / linchpin objects (ví dụ người nổi tiếng có hàng triệu follower) → một reducer phải làm nhiều hơn hẳn → cả job chờ reducer chậm nhất.
  • Skewed join (Pig): chạy job sampling để tìm hot key; record của hot key được gửi tới một reducer ngẫu nhiên trong vài reducer, còn input kia của join thì replicate tới tất cả các reducer đó.
  • Sharded join (Crunch): tương tự nhưng hot key phải chỉ định tay.
  • Hive skewed join: khai báo hot key trong metadata của bảng, lưu riêng, dùng map-side join cho chúng.
  • Grouping hot key 2 giai đoạn: stage 1 gửi tới reducer ngẫu nhiên để pre-aggregate; stage 2 gộp các kết quả trung gian lại thành một giá trị mỗi key.

Map-Side Joins

Reduce-side join: không cần giả định gì về input, nhưng sort + copy + merge rất tốn (dữ liệu có thể bị ghi ra disk nhiều lần). Nếu có thể giả định về input thì dùng map-side join: job chỉ có mapper, không reducer, không sort.

  • Broadcast hash join: một input đủ nhỏ để vừa memory của mỗi mapper. Mapper load toàn bộ user DB vào hash table rồi quét các event và lookup. "Broadcast" vì input nhỏ được gửi tới mọi partition của input lớn. Tên gọi: replicated join (Pig), MapJoin (Hive); cũng có trong Cascading, Crunch, Impala. Biến thể: lưu input nhỏ thành read-only index trên disk local (nhờ page cache mà vẫn gần nhanh như memory).
  • Partitioned hash join (bucketed map join trong Hive): cả hai input được partition cùng cách (cùng key, cùng hash function, cùng số partition) — ví dụ theo chữ số cuối của user ID. Mapper 3 chỉ load user có ID tận cùng là 3 vào hash table rồi quét event của những user đó → hash table nhỏ hơn. Thường có được khi input do các job trước tạo ra.
  • Map-side merge join: input không chỉ partition giống nhau mà còn sort theo cùng key → mapper đọc tuần tự cả hai file và merge như reducer; không cần vừa memory. Có thể tách thành job riêng nếu dataset đã sort còn dùng cho mục đích khác.

Ảnh hưởng tới workflow: output của reduce-side join được partition và sort theo join key; output của map-side join được partition và sort giống input lớn. Map-side join cần biết physical layout (số partition, key partition/sort) → metadata lưu ở HCatalog / Hive metastore.

Thuật toánĐiều kiệnCách làmƯu / nhược
Reduce-side sort-merge joinKhông giả định gìMapper emit theo join key → shuffle/sort → reducer mergeTổng quát nhất; tốn sort, network, disk I/O
Broadcast hash joinMột input nhỏ, vừa memoryMỗi mapper load input nhỏ vào hash table, quét input lớnRất nhanh, không shuffle; input nhỏ phải thật sự nhỏ
Partitioned hash joinHai input partition giống hệt nhauMỗi mapper load một partition input nhỏ, quét partition tương ứngHash table nhỏ hơn; cần layout chuẩn bị trước
Map-side merge joinPartition giống nhau và sort theo cùng keyMapper merge tuần tự hai fileKhông cần vừa memory; cần job trước đã sort sẵn

The Output of Batch Workflows

Batch không phải OLTP, cũng không hẳn analytics (output thường không phải là report) — output thường là một cấu trúc dữ liệu khác.

Building search indexes

  • Ứng dụng đầu tiên của MapReduce ở Google: build search index bằng 5–10 job. Vẫn là cách tốt để build index Lucene/Solr.
  • Mapper partition document, mỗi reducer build index (term dictionary → postings list) cho partition của mình, ghi lên HDFS. Index file immutable.
  • Document thay đổi: chạy lại toàn bộ và thay index cũ (tốn nhưng dễ suy luận: "documents in, indexes out"), hoặc build incremental (Lucene segment + merge, Ch11).

Key-value stores as batch process output

  • Output của ML (spam filter, anomaly detection, image recognition) và recommendation (people you may know, related products) thường là một database để web app query theo user ID/product ID.
  • Sai lầm: ghi thẳng từ mapper/reducer vào production DB từng record:
    • Một network request mỗi record → chậm hơn nhiều bậc so với throughput của batch.
    • Nhiều task ghi song song → làm quá tải DB, ảnh hưởng các query production.
    • Mất đảm bảo all-or-nothing: side effect bên ngoài lộ ra từ job chạy dở, retry, speculative execution.
  • Đúng: build DB file ngay trong batch job, ghi vào output directory trên HDFS, rồi bulk load vào server read-only: Voldemort, Terrapin, ElephantDB, HBase bulk loading. Mapper trích key + sort theo key đã là phần lớn việc build index; file read-only nên đơn giản, không cần WAL.
  • Voldemort: vẫn phục vụ từ file cũ trong lúc copy file mới, rồi chuyển atomic sang file mới; có lỗi thì quay lại file cũ.

Philosophy of batch process outputs

Giống Unix: input không đổi, output cũ bị thay thế hoàn toàn, không side effect. Lợi ích: - Human fault tolerance: deploy code lỗi → rollback code và chạy lại, hoặc chỉ cần trỏ về thư mục output cũ. DB read-write thì rollback code không sửa được dữ liệu đã hỏng. - Minimizing irreversibility → phát triển feature nhanh hơn (Agile). - Task lỗi được retry tự động an toàn vì input immutable và output của task lỗi bị bỏ. - Cùng input dùng cho nhiều job, kể cả job monitoring so sánh output với lần chạy trước. - Tách logic và wiring (input/output directory) → tái sử dụng: một team viết job, team khác quyết định chạy ở đâu, khi nào. - Khác Unix: Hadoop dùng format có cấu trúc (Avro, Parquet) thay vì text phải parse, hỗ trợ schema evolution.

Comparing Hadoop to Distributed Databases

Hadoop ≈ Unix phân tán: HDFS là filesystem, MapReduce là một "process Unix kỳ quặc" luôn chạy sort giữa map và reduce. Các thuật toán join song song đã có trong MPP databases (Gamma, Teradata, Tandem NonStop SQL) từ hơn 10 năm trước. Khác biệt lớn: MPP tập trung chạy SQL analytic song song; MapReduce + HDFS giống một hệ điều hành đa năng chạy chương trình tuỳ ý.

Khía cạnhMPP databasesHadoop (MapReduce + HDFS)
StoragePhải mô hình hoá dữ liệu cẩn thận trước, format độc quyềnChỉ là byte — text, ảnh, video, vector, genome…; dump dữ liệu thô trước, tính sau (data lake, schema-on-read, sushi principle: "raw data is better")
Processing modelNguyên khối, tối ưu chặt, SQL biểu đạt tốt, hợp với BI tool (Tableau)Chạy code tuỳ ý (ML, search ranking, image analysis); SQL (Hive) chỉ là một trong nhiều model; nhiều model cùng chạy trên một cluster, cùng file
Xử lý lỗiNode chết → abort cả query, chạy lại (query chỉ vài giây/phút)Retry ở mức từng task
Memory vs diskGiữ nhiều trong memory (hash join)Hăng hái ghi ra disk (vì fault tolerance và giả định dữ liệu quá lớn)
Phù hợpQuery analytic nhanhJob rất lớn, chạy lâu, trên cluster dùng chung
  • Hadoop thường làm ETL: dump dữ liệu thô từ OLTP → MapReduce làm sạch và chuyển về dạng quan hệ → import vào MPP warehouse. Data modeling vẫn diễn ra nhưng tách khỏi bước thu thập.
  • Hệ sinh thái Hadoop có cả HBase (OLTP) và Impala (MPP) — không dùng MapReduce nhưng dùng chung HDFS.

Designing for frequent faults — vì sao MapReduce "sợ lỗi" đến thế? Không phải vì phần cứng hay hỏng, mà vì môi trường mixed-use ở Google: online service và batch chạy chung máy, mỗi task có priority; task priority thấp (batch) có thể bị preempt bất cứ lúc nào để nhường tài nguyên. Task chạy 1 giờ có ~5% rủi ro bị kill (cao hơn lỗi phần cứng hơn một bậc); job 100 task mỗi task 10 phút → >50% khả năng ít nhất một task bị kill. Khả năng kill tuỳ ý cho phép overcommit tài nguyên → tận dụng máy tốt hơn. Trong môi trường ít preemption (YARN, Mesos, Kubernetes thời điểm viết sách) thì thiết kế này kém hợp lý hơn.

Beyond MapReduce

MapReduce đơn giản để hiểu nhưng khó dùng (phải tự viết join…) → có Pig, Hive, Cascading, Crunch. Nhưng mô hình thực thi của MapReduce có vấn đề về performance mà thêm abstraction cũng không sửa được. Nó rất bền bỉ (xử lý được dữ liệu cực lớn trên cluster hay kill task) nhưng chậm.

Materialization of Intermediate State

  • Publish output ra vị trí đã biết trên HDFS là hợp lý khi nhiều team dùng. Nhưng thường output chỉ là intermediate state giữa hai job của cùng một team.
  • Materialization = tính trước và ghi ra file, trái với Unix pipe stream dữ liệu qua buffer nhỏ.
  • Nhược điểm của materialization toàn phần:
    • Job sau phải chờ mọi task của job trước xong; straggler làm chậm cả workflow.
    • Mapper thường thừa: chỉ đọc lại file reducer vừa ghi để partition/sort lại.
    • Intermediate state trên HDFS bị replicate nhiều bản — quá mức cần thiết cho dữ liệu tạm.

Dataflow engines

Spark, Tez, Flink (dựa trên nghiên cứu Dryad, Nephele): coi cả workflow là một job, mô hình hoá luồng dữ liệu qua các operator (không bắt buộc xen kẽ map/reduce). Các cách nối operator: repartition + sort (như shuffle); partition nhưng không sort (cho partitioned hash join); broadcast (cho broadcast hash join).

Khía cạnhMapReduceDataflow engines (Spark, Flink, Tez)
Đơn vịNhiều job độc lập nối qua HDFSCả workflow là một job (DAG operator)
SortLuôn sort giữa map và reduceChỉ sort khi thật sự cần
Mapper thừaCóKhông — gộp vào operator trước
Intermediate stateGhi HDFS, replicateMemory hoặc disk local
Bắt đầu stage sauChờ stage trước xong hoàn toànOperator chạy ngay khi input sẵn sàng (pipelined)
LocalityHạn chếScheduler biết toàn bộ dependency → đặt producer/consumer cùng máy, trao đổi qua shared memory
JVMKhởi động JVM mới cho mỗi taskTái sử dụng JVM
Fault toleranceĐọc lại input bền vững từ HDFSRecompute từ lineage (RDD) hoặc checkpoint (Flink)

Workflow viết bằng Pig/Hive/Cascading có thể chuyển từ MapReduce sang Tez/Spark chỉ bằng đổi config. Tez là library mỏng dùng YARN shuffle service; Spark và Flink là framework lớn có network layer, scheduler và API riêng.

Fault tolerance

  • Không materialize thì khi mất máy phải recompute từ stage trước (nếu còn) hoặc từ input gốc trên HDFS. Spark dùng RDD để theo dõi lineage (dữ liệu được tính từ partition nào, qua operator nào); Flink checkpoint state của operator.
  • Cần determinism: nếu dữ liệu tính lại khác dữ liệu đã gửi xuống downstream thì phải kill và chạy lại cả downstream. Nguồn nondeterminism dễ lọt vào: thứ tự duyệt hash table, số ngẫu nhiên (dùng seed cố định), system clock, dữ liệu ngoài.
  • Nếu intermediate data nhỏ hơn nhiều so với input, hoặc tính toán rất nặng CPU → materialize lại rẻ hơn recompute.

Discussion of materialization

MapReduce ≈ ghi output mỗi lệnh ra file tạm; dataflow engines ≈ Unix pipes (Flink đặc biệt hướng tới pipelined execution). Sort vẫn buộc phải tiêu thụ hết input trước khi ra output. Input và output cuối vẫn nằm trên HDFS, immutable, output được thay thế hoàn toàn — chỉ bỏ được phần intermediate.

Graphs and Iterative Processing

  • Xử lý batch trên toàn bộ graph (khác OLTP graph query ở Ch2): ví dụ PageRank, recommendation.
  • Đừng nhầm: DAG trong dataflow engine là luồng dữ liệu có dạng graph; graph processing thì bản thân dữ liệu là graph.
  • Nhiều thuật toán graph = lan truyền thông tin theo cạnh, lặp tới khi hội tụ (ví dụ transitive closure). MapReduce chỉ đi một lượt, nên phải có scheduler bên ngoài lặp lại: chạy một bước → kiểm tra điều kiện dừng → lặp. Rất kém hiệu quả vì mỗi vòng đọc lại toàn bộ input và ghi toàn bộ output mới, dù chỉ một phần nhỏ graph thay đổi.

The Pregel processing model

  • Bulk synchronous parallel (BSP), phổ biến nhờ paper Pregel của Google; có trong Apache Giraph, Spark GraphX, Flink Gelly.
  • Giống mapper "gửi message" tới reducer: ở đây một vertex gửi message tới vertex khác (thường dọc theo cạnh).
  • Mỗi iteration (superstep), một function được gọi cho mỗi vertex với mọi message gửi tới nó. Khác MapReduce: vertex nhớ state trong memory giữa các iteration → chỉ xử lý message mới; phần graph không có message thì không tốn công.
  • Giống actor model, nhưng state và message bền vững, chịu lỗi, và giao tiếp theo vòng cố định.

Fault tolerance (Pregel)

  • Chỉ giao tiếp qua message → batch được, ít chờ; chỉ chờ ở ranh giới giữa các iteration.
  • Đảm bảo message được xử lý đúng một lần ở iteration sau, dù mạng mất/trùng/trễ.
  • Checkpoint state của mọi vertex ở cuối iteration; lỗi thì rollback cả graph về checkpoint gần nhất (hoặc chỉ phục hồi partition bị mất nếu thuật toán deterministic và message được log).

Parallel execution

  • "Thinking like a vertex": vertex chỉ gửi tới vertex ID, framework quyết định partition và routing.
  • Partition tối ưu rất khó → thường partition theo vertex ID tuỳ ý → overhead giao tiếp giữa các máy lớn, message trung gian có thể lớn hơn cả graph gốc.
  • Graph vừa memory một máy → thuật toán single-machine (thậm chí single-thread) thường nhanh hơn; vừa disk một máy → GraphChi. Chỉ dùng Pregel phân tán khi thực sự không vừa một máy.

High-Level APIs and Languages

  • Hạ tầng đã trưởng thành (hàng PB, cluster >10.000 máy) → trọng tâm chuyển sang programming model và hiệu quả.
  • Hive, Pig, Cascading, Crunch và API của Spark/Flink (lấy cảm hứng từ FlumeJava) dùng các khối kiểu quan hệ: join, group by, filter, aggregate. Lợi ích: ít code hơn, dùng interactive trong shell (giống tinh thần Unix), và chạy hiệu quả hơn ở mức máy.

The move toward declarative query languages

  • Khai báo join theo kiểu quan hệ → cost-based optimizer (Hive, Spark, Flink) tự chọn thuật toán join phù hợp và đổi thứ tự join để giảm intermediate state. Không cần nhớ hết các loại join.
  • Nhưng vẫn giữ lợi thế của callback tuỳ ý: dùng được hệ sinh thái thư viện (parsing, NLP, image, numeric) — UDF trong DB thì thường cồng kềnh, không tích hợp với Maven/npm/Rubygems.
  • Filter/map đơn giản viết declarative → optimizer tận dụng column-oriented storage (chỉ đọc cột cần), vectorized execution; Spark sinh JVM bytecode, Impala dùng LLVM sinh native code.
  • → Batch framework ngày càng giống MPP (performance tương đương) nhưng vẫn giữ tính linh hoạt.

Specialization for different domains

  • ML: Mahout (trên MapReduce/Spark/Flink), MADlib (trong MPP, Apache HAWQ).
  • Spatial: k-nearest neighbors (similarity search); genome analysis cần approximate string matching.
  • Kết luận: batch engine và MPP database đang hội tụ — rốt cuộc đều là hệ thống lưu trữ và xử lý dữ liệu.

Summary

  • Hai bài toán cốt lõi của distributed batch: partitioning (gom dữ liệu liên quan về cùng chỗ; dataflow engines tránh sort khi không cần) và fault tolerance (MapReduce ghi disk nhiều nên phục hồi task dễ nhưng chậm khi không có lỗi; dataflow engines giữ trong memory, recompute nhiều hơn khi lỗi; operator deterministic giảm lượng cần tính lại).
  • Programming model cố ý bị giới hạn: callback stateless, không side effect ngoài output → framework che giấu các vấn đề phân tán; nếu nhiều task của cùng partition thành công thì chỉ một output được hiển thị. Output cuối như thể không có lỗi — mạnh hơn nhiều so với online service ghi DB như side effect.
  • Batch = input bounded, output dẫn xuất từ input, job có điểm kết thúc. Stream (Ch11) = input unbounded, job không bao giờ kết thúc.

⚠️ Hiểu lầm & cạm bẫy thường gặp

  • "Shuffle" nghĩa là xáo trộn ngẫu nhiên: sai — đó là partition theo hash, sort và copy, hoàn toàn deterministic.
  • Gọi DB/service từ xa trong mapper cho mỗi record: chậm, làm quá tải DB, và làm job nondeterministic. Hãy copy dữ liệu vào HDFS rồi join.
  • Ghi thẳng output vào production DB từ batch job: mất all-or-nothing, dễ ghi trùng khi retry/speculative execution. Hãy build file rồi bulk load và swap atomic.
  • Bỏ qua skew: một hot key làm cả job chờ một reducer. Cần sampling + skewed join hoặc aggregate 2 giai đoạn.
  • Dùng broadcast hash join khi bảng "nhỏ" không thật sự nhỏ → OOM trên mapper. (Spark có ngưỡng autoBroadcastJoinThreshold.)
  • Partitioned hash join khi hai input partition khác nhau (khác số partition hoặc hash function) → join sai hoặc thiếu kết quả.
  • Nghĩ MapReduce chịu lỗi tốt vì phần cứng hay hỏng: lý do thật là preemption trong cluster mixed-use.
  • Nhầm DAG của dataflow engine với graph processing.
  • Operator nondeterministic (thứ tự duyệt hash map, random không seed, now()) trong Spark/Flink → recompute sinh ra kết quả khác, gây lỗi dây chuyền.
  • Mặc định dùng hệ phân tán cho graph: graph vừa một máy thì single-machine thường nhanh hơn nhiều.
  • Nghĩ MapReduce là thứ mới: MPP database đã làm join song song từ thập niên 80–90; điểm mới là tính đa năng và phần cứng phổ thông.

💼 Áp dụng thực tế & phỏng vấn

  • Hệ thống hiện đại: Spark (Databricks, EMR), Flink, Hive/Trino/Presto trên data lake (S3 + Parquet/ORC, Iceberg/Delta/Hudi), orchestration bằng Airflow/Dagster. BigQuery, Snowflake là hậu duệ kiểu MPP nhưng tách storage và compute.
  • Thiết kế pipeline ETL/ELT: raw zone (dữ liệu thô, immutable) → cleaned → curated; mỗi job ghi vào một partition/thư mục mới rồi chuyển con trỏ (write-audit-publish) → rollback dễ dàng. Chính là triết lý "input immutable, output thay thế".
  • Recommendation / "People You May Know": batch job hằng đêm tính kết quả → bulk load vào KV store read-only (kiểu Voldemort/Terrapin; ngày nay là Cassandra/HBase bulk load, RocksDB SST ingest) → web app đọc với latency thấp. Hay gặp trong phỏng vấn thiết kế "feed", "recommendation", "top-K".
  • Chọn join strategy trong Spark: broadcast join cho dimension nhỏ; sort-merge join mặc định cho hai bảng lớn; bucketing theo join key để tránh shuffle; AQE xử lý skew (tách partition lệch). Giải thích được vì sao shuffle đắt (network + disk + sort).
  • Top-K / word count / sessionization là câu hỏi kinh điển: trình bày được map (emit key), shuffle, reduce, và combiner/pre-aggregate để giảm dữ liệu shuffle; aggregate 2 giai đoạn cho hot key.
  • Lambda/Kappa architecture: batch cho độ chính xác và reprocessing, stream cho độ tươi — nối sang Ch11.
  • Backfill và reprocessing: nhờ input immutable, sửa bug xong chạy lại trên dữ liệu lịch sử. Nhớ nhắc tới yêu cầu determinism và idempotent output.
  • Ước lượng: batch được tối ưu cho throughput (MB/s/node × số node); đọc tuần tự tốt hơn random I/O; data locality giảm network.

❓ Câu hỏi ôn tập

Bấm vào câu hỏi để xem đáp án
1Vì sao pipeline Unix sort | uniq -c scale tốt hơn hash table trong memory khi dữ liệu lớn?
Hash table cần working set (số key khác nhau) vừa memory. Sort thì sort từng chunk trong memory, spill ra disk rồi merge tuần tự (giống SSTable), nên dữ liệu lớn hơn RAM vẫn chạy; GNU sort còn song song hoá trên nhiều core.
2Mô tả luồng thực thi một MapReduce job, bao gồm shuffle.
Mỗi block input → một map task (chạy gần dữ liệu). Mapper emit key-value, chia theo hash(key) cho từng reducer, sort và ghi ra disk local. Reducer kéo partition của mình từ mọi mapper (shuffle), merge giữ thứ tự, gọi reduce cho từng key với iterator các value, ghi output lên HDFS.
3Sort-merge join (reduce-side) hoạt động thế nào trong ví dụ user activity? Secondary sort dùng để làm gì?
Hai nhóm mapper emit theo userID (event, và ngày sinh từ user DB). Shuffle gom cùng userID về một reducer; secondary sort đảm bảo record user đứng trước các event, nên reducer chỉ cần giữ ngày sinh trong một biến rồi duyệt event, emit (url, tuổi). Không cần gọi network, dùng rất ít memory.
4So sánh broadcast hash join, partitioned hash join và map-side merge join.
Broadcast: một input nhỏ vừa memory, load vào mọi mapper. Partitioned: hai input partition giống hệt nhau, mỗi mapper load một partition. Merge join: partition giống nhau và sort theo cùng key, mapper merge tuần tự nên không cần vừa memory. Cả ba đều không có reducer/sort; output được partition giống input lớn.
5Xử lý hot key (skew) trong join và group by thế nào?
Join: sampling để tìm hot key, rải record của hot key sang vài reducer ngẫu nhiên và replicate phía input kia tới các reducer đó (Pig skewed join, Crunch sharded join), hoặc dùng map-side join cho hot key (Hive). Group by: aggregate 2 giai đoạn — stage 1 gửi tới reducer ngẫu nhiên để pre-aggregate, stage 2 gộp lại.
6Vì sao không nên ghi trực tiếp từ batch job vào production database? Cách làm tốt hơn?
Mỗi record một request thì chậm, nhiều task ghi song song làm quá tải DB, và mất đảm bảo all-or-nothing (side effect lộ ra khi retry hoặc job fail). Nên build file DB trong job, ghi lên HDFS, rồi bulk load vào store read-only và chuyển atomic (Voldemort, HBase bulk load).
7Vì sao MapReduce ghi mọi thứ ra disk và retry theo từng task? So sánh với MPP.
Được thiết kế cho cluster mixed-use của Google, nơi task batch priority thấp hay bị preempt (~5%/giờ). Retry mức task tránh phải chạy lại cả job lớn. MPP thì abort cả query khi có lỗi (query ngắn) và giữ dữ liệu trong memory để nhanh.
8Dataflow engines cải thiện gì so với MapReduce và chịu lỗi ra sao? Pregel khác gì?
Coi cả workflow là một DAG operator: chỉ sort khi cần, bỏ mapper thừa, intermediate state để trong memory/disk local, pipeline giữa các stage, tái sử dụng JVM, locality tốt hơn. Chịu lỗi bằng recompute theo lineage (Spark RDD) hoặc checkpoint (Flink), nên cần operator deterministic. Pregel (BSP) dành cho graph lặp: vertex giữ state giữa các iteration, gửi message dọc theo cạnh, checkpoint cuối mỗi iteration.

Đây là bản tóm tắt và ghi chú, không thay thế sách gốc. Hãy ủng hộ tác giả bằng cách đọc bản gốc.