Chương 11 — Stream Processing
Phần III — Derived Data

Chương 11 — Stream Processing

20 phút đọc DDIA · Martin Kleppmann

🎯 Mục tiêu chương: Hiểu cách biểu diễn, truyền tải và xử lý event stream (dữ liệu unbounded, đến dần theo thời gian) — từ message broker truyền thống đến log-based broker như Kafka; thấy được mối liên hệ sâu sắc giữa database và stream (CDC, event sourcing); và nắm các vấn đề khó khi xử lý stream: thời gian, window, join, fault tolerance / exactly-once.

Chương 10 giả định input là bounded (kích thước hữu hạn, biết khi nào đọc xong). Thực tế dữ liệu hầu hết là unbounded: user sinh dữ liệu hôm qua, hôm nay, và cả ngày mai. Batch phải cắt dữ liệu thành các lát cố định (mỗi ngày, mỗi giờ) → output trễ. Stream processing đẩy ý tưởng này tới cực hạn: bỏ lát thời gian, xử lý từng event ngay khi nó xảy ra.

Transmitting Event Streams

  • Event = tương đương "record" trong batch: object nhỏ, tự chứa, immutable, mô tả một việc đã xảy ra tại một thời điểm (thường kèm timestamp theo time-of-day clock). Ví dụ: user xem trang, mua hàng, sensor đo nhiệt độ, metric CPU.
  • Event được encode (JSON, Avro, binary…) để lưu hoặc gửi qua mạng.
  • Producer (publisher/sender) sinh event một lần; nhiều consumer (subscriber/recipient) có thể xử lý. Các event liên quan được nhóm vào một topic hay stream (tương tự filename trong batch).
  • Về nguyên tắc có thể dùng file/DB: producer ghi, consumer poll định kỳ. Nhưng poll càng dày thì càng nhiều request trả về rỗng → overhead lớn. Tốt hơn là consumer được notify khi có event mới. DB truyền thống chỉ có trigger — khá hạn chế, là "tính năng thêm vào sau" → cần công cụ chuyên dụng.

Messaging Systems

Mô hình publish/subscribe: nhiều producer gửi vào cùng topic, nhiều consumer nhận. Để phân biệt các hệ thống, hỏi 2 câu:

  1. Producer gửi nhanh hơn consumer xử lý thì sao? Có 3 lựa chọn:
    • Drop message.
    • Buffer trong queue — cần biết queue phình to thì sao: crash khi hết RAM hay spill ra disk (và disk ảnh hưởng hiệu năng thế nào)?
    • Backpressure (flow control): chặn producer lại. Unix pipe và TCP dùng cách này (buffer nhỏ cố định, đầy thì sender bị block).
  2. Node crash / offline thì có mất message không? Durability cần ghi disk và/hoặc replication — có chi phí. Nếu chấp nhận mất đôi chút thì throughput cao hơn, latency thấp hơn. Sensor metric định kỳ mất 1 điểm không sao; nhưng đếm event thì mỗi message mất là counter sai.

Direct messaging từ producer tới consumer

Không qua node trung gian: - UDP multicast — phổ biến trong tài chính (stock feed) vì latency thấp; tầng ứng dụng tự retransmit gói mất. - Thư viện brokerless: ZeroMQ, nanomsg (pub/sub qua TCP/IP multicast). - StatsD, Brubeck dùng UDP không tin cậy để thu metric → counter chỉ gần đúng. - Webhooks: đăng ký callback URL, service kia gọi HTTP khi có event.

Hạn chế: ứng dụng phải tự lo khả năng mất message; thường giả định producer và consumer luôn online. Consumer offline sẽ mất message; producer crash thì mất buffer retry.

Message brokers

Message broker (message queue) = một kiểu database tối ưu cho message stream, chạy như server; producer và consumer là client. - Tập trung dữ liệu ở broker → chịu được client connect/disconnect/crash; durability chuyển về broker (tùy cấu hình in-memory hay ghi disk). - Với consumer chậm, thường cho phép unbounded queueing. - Consumer thường asynchronous: producer chỉ đợi broker xác nhận đã buffer, không đợi consumer xử lý.

So sánh message broker với database

  • DB giữ dữ liệu tới khi xóa tường minh; broker (kiểu truyền thống) tự xóa message khi đã giao thành công → không hợp lưu trữ lâu dài.
  • Broker giả định working set nhỏ (queue ngắn); nếu phải buffer nhiều (spill ra disk) thì throughput giảm.
  • DB có secondary index, query tùy ý; broker chỉ hỗ trợ subscribe theo pattern topic.
  • Query DB là snapshot tại một thời điểm, không biết khi kết quả lỗi thời; broker notify client khi có dữ liệu mới nhưng không cho query tùy ý.

Đây là góc nhìn truyền thống theo chuẩn JMS và AMQP: RabbitMQ, ActiveMQ, HornetQ, Qpid, TIBCO EMS, IBM MQ, Azure Service Bus, Google Cloud Pub/Sub.

Multiple consumers

  • Load balancing: mỗi message giao cho một consumer → chia việc, hợp khi message tốn công xử lý (AMQP: nhiều client cùng consume một queue; JMS: shared subscription).
  • Fan-out: mỗi message giao cho tất cả consumer → giống nhiều batch job cùng đọc một file input (JMS topic subscription; AMQP exchange binding).
  • Kết hợp được: nhiều consumer group, mỗi group nhận toàn bộ message, trong group thì chia tải.

Acknowledgments and redelivery

  • Consumer có thể crash giữa chừng → broker dùng acknowledgment: client phải báo tường minh đã xử lý xong, broker mới xóa. Mất kết nối / timeout mà chưa ack → broker giao lại cho consumer khác.
  • Lưu ý: có thể message đã xử lý xong nhưng ack bị mất trên mạng → xử lý 2 lần; muốn tránh cần atomic commit.
  • Load balancing + redelivery ⇒ reordering: consumer 2 crash khi xử lý m3, cùng lúc consumer 1 xử lý m4; m3 được giao lại cho consumer 1 → thứ tự thành m4, m3, m5. Dù JMS/AMQP cố giữ thứ tự, sự kết hợp này tất yếu gây đảo thứ tự. Cách tránh: mỗi consumer một queue riêng. Quan trọng khi message có causal dependency.

Partitioned Logs

Khác biệt tư duy: gửi packet/message là transient, không để lại dấu vết; còn DB/file là permanent. Với broker AMQP/JMS, nhận message là destructive (ack → xóa) nên không chạy lại consumer để có cùng kết quả; consumer mới chỉ nhận message từ lúc đăng ký. Ngược lại batch job có thể chạy lại thoải mái vì input read-only.

→ Log-based message broker: kết hợp durable storage của DB với notification độ trễ thấp của messaging.

Dùng log để lưu message

  • Log = chuỗi record append-only trên disk (đã gặp ở Ch3 log-structured storage, WAL; Ch5 replication log).
  • Producer append vào cuối log; consumer đọc tuần tự; tới cuối thì chờ notify (giống tail -f).
  • Để scale vượt 1 disk: partition log, mỗi partition trên máy khác nhau, đọc/ghi độc lập. Topic = nhóm partition cùng loại message.
  • Mỗi partition gán offset tăng đơn điệu cho từng message → trong partition có total order; giữa các partition không có đảm bảo thứ tự.
  • Ví dụ: Apache Kafka, Amazon Kinesis Streams, Twitter DistributedLog. Dù ghi mọi thứ ra disk vẫn đạt hàng triệu message/giây nhờ partitioning, và fault tolerance nhờ replication.

So sánh log với messaging truyền thống

  • Fan-out là tự nhiên: đọc không xóa, nhiều consumer đọc độc lập.
  • Load balancing: broker gán nguyên partition cho từng node trong consumer group; node đó đọc tuần tự, single-thread. Nhược điểm:
    • Số consumer tối đa = số partition của topic.
    • Một message chậm chặn các message sau trong cùng partition (head-of-line blocking).
  • Kết luận: message tốn công, cần song song theo từng message, không cần thứ tự → JMS/AMQP. Throughput cao, mỗi message xử lý nhanh, thứ tự quan trọng → log-based.

Consumer offsets

  • Đọc tuần tự → chỉ cần nhớ offset hiện tại: mọi message < offset đã xử lý. Broker không cần track ack từng message, chỉ định kỳ ghi offset → ít bookkeeping, dễ batching/pipelining → throughput cao.
  • Offset rất giống log sequence number trong single-leader replication: broker ≈ leader, consumer ≈ follower.
  • Consumer chết → node khác nhận partition, đọc từ offset đã ghi cuối cùng; các message đã xử lý nhưng chưa commit offset sẽ bị xử lý lại (at-least-once).

Disk space usage

  • Log được chia segment; segment cũ bị xóa hoặc archive → log thực chất là circular/ring buffer trên disk, kích thước lớn.
  • Tính nhanh: disk 6 TB, ghi tuần tự 150 MB/s → ~11 giờ mới đầy ở tốc độ tối đa; thực tế không dùng hết băng thông nên giữ được vài ngày đến vài tuần message.
  • Throughput log gần như không đổi bất kể giữ bao lâu (luôn ghi disk), khác broker in-memory: nhanh khi queue ngắn, chậm hẳn khi phải spill ra disk.

Khi consumer không theo kịp producer

  • Log-based = buffering với buffer lớn nhưng cố định. Consumer tụt quá xa (offset trỏ tới segment đã xóa) sẽ mất message.
  • Có thể monitor consumer lag và alert; buffer lớn cho con người đủ thời gian sửa.
  • Chỉ consumer chậm bị ảnh hưởng, không làm phiền consumer khác → có thể đọc log production để dev/test/debug mà không sợ ảnh hưởng. Consumer tắt thì chỉ còn lại offset, không tốn tài nguyên (khác broker truyền thống: phải nhớ xóa queue của consumer đã tắt, không thì queue tích tụ ăn RAM).

Replaying old messages

  • Consume trong log-based broker là read-only, chỉ làm offset tiến lên, và offset do consumer kiểm soát → có thể chạy bản sao consumer với offset hôm qua, ghi output sang chỗ khác, để reprocess — lặp lại bao nhiêu lần cũng được với code khác nhau.
  • Làm messaging giống batch: derived data tách khỏi input bằng một phép biến đổi lặp lại được → dễ thử nghiệm, dễ phục hồi khi có bug.
Tiêu chíAMQP/JMS-style brokerLog-based broker (Kafka, Kinesis)
Sau khi xử lýMessage bị xóa khi ackMessage vẫn còn trong log tới khi segment hết hạn
Load balancingTheo từng message, số consumer tùy ýTheo partition, số consumer ≤ số partition
Thứ tựCó thể đảo khi redelivery + load balancingTotal order trong partition
Theo dõi tiến độAck từng messageCommit consumer offset định kỳ
ReplayKhông (đọc là destructive)Có, reset offset để đọc lại
Consumer mớiChỉ nhận message sau khi đăng kýCó thể đọc từ đầu log
Message chậmKhông chặn message khácHead-of-line blocking trong partition
Phù hợpTask queue, async RPC, job tốn công, không cần thứ tựThroughput cao, cần thứ tự, stream processing, derived data

Databases and Streams

Log-based broker lấy ý tưởng từ DB đưa vào messaging; chiều ngược lại cũng được: đưa ý tưởng stream vào DB. Một write vào database chính là một event. Replication log là stream các write event; nguyên lý state machine replication: mọi replica xử lý cùng các event theo cùng thứ tự (deterministic) → cùng trạng thái cuối. Tất cả đều là event stream.

Keeping Systems in Sync

Ứng dụng thực tế kết hợp nhiều hệ thống: OLTP DB, cache, full-text index, data warehouse — mỗi cái giữ bản sao dữ liệu theo cách riêng, cần đồng bộ. - Data warehouse thường dùng ETL batch (dump toàn bộ DB, transform, bulk-load). - Nếu dump quá chậm, người ta dùng dual writes: code ứng dụng tự ghi vào DB, rồi update search index, rồi invalidate cache…

Vấn đề của dual writes: 1. Race condition (Figure 11-4): client 1 set X=A, client 2 set X=B, đồng thời. DB nhận A rồi B → B; search index nhận B rồi A → A. Hai hệ thống vĩnh viễn lệch nhau dù không có lỗi nào xảy ra, và không ai phát hiện nếu không có cơ chế như version vector. 2. Partial failure: ghi DB thành công, ghi index thất bại → không nhất quán. Muốn "cả hai hoặc không cái nào" là bài toán atomic commit — đắt (2PC).

Gốc rễ: không có một leader duy nhất quyết định thứ tự ghi — DB có leader, index có leader, không ai follow ai (giống multi-leader conflict). Giải pháp lý tưởng: DB là leader duy nhất, biến search index thành follower của DB.

Change Data Capture

CDC = quan sát mọi thay đổi ghi vào DB và trích xuất ra dạng có thể replicate sang hệ thống khác — đặc biệt hữu ích khi thay đổi được cung cấp dưới dạng stream, ngay khi ghi. Replication log vốn bị coi là chi tiết nội bộ, không phải public API, nên trước đây rất khó làm.

  • Áp dụng log thay đổi theo đúng thứ tự → search index khớp với DB. Search index, cache, data warehouse đều chỉ là consumer của change stream (Figure 11-5).

Implementing change data capture

  • Các consumer là derived data system; DB nguồn là system of record. CDC biến DB nguồn thành leader, các hệ thống khác thành follower.
  • Log-based message broker rất hợp để vận chuyển change event vì giữ thứ tự (tránh reorder như Figure 11-2).
  • Cách implement:
    • Trigger ghi vào bảng changelog — mong manh, overhead lớn.
    • Parse replication log — robust hơn nhưng phải xử lý schema change.
  • Công cụ: LinkedIn Databus, Facebook Wormhole, Yahoo Sherpa; Bottled Water (PostgreSQL, decode WAL); Maxwell, Debezium (MySQL binlog); Mongoriver (MongoDB oplog); GoldenGate (Oracle).
  • CDC thường asynchronous: DB không chờ consumer → thêm consumer chậm không ảnh hưởng nhiều, nhưng gặp mọi vấn đề replication lag.

Initial snapshot

  • Có toàn bộ log từ đầu thì replay được toàn bộ DB, nhưng giữ mãi thì tốn disk, replay thì lâu → log bị cắt.
  • Dựng index mới cần bản sao đầy đủ → bắt đầu từ consistent snapshot, và snapshot phải ứng với một offset xác định trong change log để biết áp thay đổi từ đâu (giống setup follower mới).

Log compaction

  • Giống ở Ch3: giữ bản ghi mới nhất cho mỗi key, bỏ các bản cũ; tombstone (null) nghĩa là xóa key.
  • Dung lượng log compacted chỉ phụ thuộc nội dung hiện tại của DB, không phụ thuộc số lần ghi.
  • Với CDC: mỗi change có primary key, update thay thế giá trị cũ → dựng derived system mới chỉ cần consumer mới đọc từ offset 0 của topic compacted — không cần snapshot DB nguồn nữa.
  • Kafka hỗ trợ log compaction → broker dùng được cho durable storage, không chỉ transient messaging.

API support for change streams

DB bắt đầu hỗ trợ change stream như first-class interface: RethinkDB (subscribe khi kết quả query đổi), Firebase, CouchDB (change feed), Meteor (dùng MongoDB oplog để cập nhật UI), VoltDB (output stream là một bảng chỉ insert, không query được). Kafka Connect tích hợp CDC cho nhiều DB vào Kafka.

Event Sourcing

Kỹ thuật từ cộng đồng Domain-Driven Design (DDD): lưu mọi thay đổi trạng thái ứng dụng như log các change event — giống CDC nhưng ở mức trừu tượng khác:

  • CDC: ứng dụng dùng DB kiểu mutable (update/delete tùy ý); log được trích ở mức thấp (replication log), đảm bảo đúng thứ tự ghi thật; ứng dụng không cần biết có CDC.
  • Event sourcing: logic ứng dụng xây trực tiếp trên immutable event trong event log; event store append-only, update/delete bị hạn chế hoặc cấm; event phản ánh việc xảy ra ở mức ứng dụng, không phải thay đổi state mức thấp.

Ví dụ: event "sinh viên hủy đăng ký môn học" diễn tả rõ ý định, trung lập; còn side effect "xóa một dòng trong bảng enrollments, thêm lý do hủy vào bảng feedback" gắn nhiều giả định về cách dữ liệu sẽ được dùng. Khi có tính năng mới "nhường chỗ cho người tiếp theo trong waiting list", event sourcing chỉ cần móc thêm side effect mới vào event có sẵn.

Lợi ích: dễ tiến hóa ứng dụng, dễ debug (hiểu tại sao việc gì đã xảy ra), chống bug ứng dụng. Liên quan: chronicle data model, fact table trong star schema. Có DB chuyên dụng (Event Store) nhưng cũng làm được với DB thường hay log-based broker.

Deriving current state from the event log

  • User cần state hiện tại (giỏ hàng hiện có gì), không phải lịch sử thay đổi → ứng dụng phải biến log (dạng write) thành state (dạng read) bằng logic deterministic để chạy lại được.
  • Log compaction khác nhau:
    • CDC: event chứa toàn bộ bản mới của record → giữ event mới nhất của key là đủ → compact được.
    • Event sourcing: event ở mức ý định, event sau không ghi đè event trước → cần toàn bộ lịch sử → không compact theo cách đó được.
  • Thường lưu snapshot state để khỏi replay toàn bộ — chỉ là tối ưu hiệu năng; ý định vẫn là giữ mọi raw event mãi mãi.

Commands and events

  • Request từ user ban đầu là command — có thể thất bại (vi phạm ràng buộc). Validate xong và được chấp nhận → trở thành event: durable, immutable, là một fact.
  • Ví dụ: đăng ký username, đặt ghế máy bay/rạp hát → phải kiểm tra còn trống trước. Nếu sau đó khách hủy, việc họ đã từng đặt vẫn là sự thật; hủy là một event riêng thêm vào sau.
  • Consumer không được reject event (đã thành một phần immutable của log, consumer khác có thể đã thấy) → validation phải xảy ra đồng bộ trước khi thành event (ví dụ serializable transaction vừa validate vừa publish).
  • Hoặc tách làm 2 event: tentative reservation rồi confirmation sau khi validate → cho phép validate bất đồng bộ.
Tiêu chíChange Data CaptureEvent Sourcing
Mức trừu tượngThấp: thay đổi row/documentCao: hành động/ý định ở mức nghiệp vụ
Cách ứng dụng dùng DBMutable, update/delete tự doAppend-only immutable event
Ứng dụng có biết không?Không cần biếtĐược thiết kế xoay quanh event log
Nguồn logTrích từ replication log / WAL / binlogEvent store do ứng dụng ghi
Log compactionĐược (event mới nhất của key chứa đủ state)Không theo cách đó (cần toàn bộ lịch sử); dùng snapshot
Công cụDebezium, Maxwell, Databus, GoldenGate, Kafka ConnectEvent Store, Kafka, DB thường

State, Streams, and Immutability

  • Mọi state thay đổi đều là kết quả của chuỗi event đã tác động lên nó: danh sách ghế trống = kết quả các lượt đặt; số dư = tổng credit/debit; biểu đồ response time = tổng hợp từng request.
  • Mutable state và append-only log immutable không mâu thuẫn — là hai mặt của một đồng xu. Changelog biểu diễn sự tiến hóa của state theo thời gian.
  • Góc nhìn toán học (Figure 11-6): state = tích phân event stream theo thời gian; change stream = đạo hàm của state.
  • Lưu changelog bền vững → state tái tạo được. Coi log là system of record, state là derived → dễ suy luận về dataflow. Pat Helland: "The truth is the log. The database is a cache of a subset of the log."
  • Log compaction là cầu nối: chỉ giữ phiên bản mới nhất của mỗi record.

Advantages of immutable events

  • Kế toán dùng immutability hàng thế kỷ: ledger append-only; sai thì thêm giao dịch bù trừ (refund) chứ không xóa/sửa. Giao dịch sai vẫn nằm đó cho audit.
  • Deploy code lỗi ghi dữ liệu sai: nếu code có thể ghi đè phá hủy dữ liệu thì phục hồi rất khó; với log immutable thì dễ chẩn đoán và phục hồi.
  • Event chứa nhiều thông tin hơn state: khách thêm món vào giỏ rồi bỏ ra — với fulfillment thì triệt tiêu nhau, nhưng với analytics thì biết khách đã từng cân nhắc món đó. DB xóa item sẽ mất thông tin này.

Deriving several views from the same event log

  • Tách mutable state khỏi immutable log → dựng nhiều read view từ cùng log: Druid ingest trực tiếp từ Kafka, Pistachio dùng Kafka làm commit log, Kafka Connect sink export sang nhiều DB/index.
  • Tính năng mới cần trình bày dữ liệu kiểu mới → dựng view mới song song từ log, không cần sửa hệ thống cũ; chạy song song dễ hơn schema migration phức tạp; xong thì tắt view cũ.
  • CQRS (Command Query Responsibility Segregation): tách dạng dữ liệu khi ghi khỏi dạng khi đọc. Sai lầm truyền thống là nghĩ dữ liệu phải ghi theo đúng dạng sẽ query. Tranh luận normalize vs denormalize trở nên ít quan trọng: denormalize trong read view là hợp lý vì quá trình dịch từ log giữ nó nhất quán.
  • Ví dụ Twitter home timeline: cache denormalize cao (tweet được nhân bản vào timeline của mọi follower), fan-out service giữ nó đồng bộ.

Concurrency control

  • Nhược điểm lớn nhất: consumer asynchronous → user ghi rồi đọc view có thể chưa thấy write của mình (read-your-own-writes). Giải pháp: cập nhật view đồng bộ với append log (cần transaction cùng storage hoặc distributed transaction), hoặc dùng total order broadcast để làm linearizable storage.
  • Mặt khác, event sourcing đơn giản hóa concurrency: nhiều multi-object transaction chỉ cần vì một hành động phải sửa nhiều chỗ; nếu thiết kế event là mô tả tự chứa của hành động thì chỉ cần một write duy nhất (append) — dễ atomic.
  • Nếu log và state được partition giống nhau, consumer single-thread mỗi partition không cần concurrency control — log định nghĩa thứ tự tuần tự (giống actual serial execution).

Limitations of immutability

  • Nhiều hệ thống dùng immutability nội bộ: MVCC/snapshot, Git/Mercurial/Fossil.
  • Khả thi hay không tùy churn: workload chủ yếu thêm dữ liệu thì dễ; workload update/delete nhiều trên dataset nhỏ thì lịch sử phình to, fragmentation, compaction/GC trở thành vấn đề vận hành.
  • Có lúc phải xóa thật vì lý do pháp lý: privacy regulation (xóa dữ liệu cá nhân khi đóng tài khoản), dữ liệu sai, rò rỉ dữ liệu nhạy cảm. Thêm event "đã xóa" là không đủ — cần viết lại lịch sử: Datomic gọi là excision, Fossil gọi là shunning.
  • Xóa thật rất khó: bản sao nằm ở storage engine, filesystem, SSD (ghi chỗ mới thay vì ghi đè), backup (cố ý immutable). Xóa thường chỉ là "làm khó lấy lại" chứ không phải "không thể lấy lại".

Processing Streams

Có stream rồi thì làm gì? Ba lựa chọn: 1. Ghi vào DB/cache/search index để client khác query (tương đương output của batch workflow). 2. Đẩy tới người dùng: email, push notification, real-time dashboard — con người là consumer cuối. 3. Xử lý một/nhiều input stream để tạo output stream — qua pipeline nhiều stage. Code như vậy gọi là operator hay job.

Stream processor đọc input read-only, ghi output append-only — giống Unix process/MapReduce; partitioning, parallelization, map/filter tương tự. Khác biệt cốt lõi: stream không bao giờ kết thúc → không sort được (không dùng sort-merge join), và fault tolerance phải khác (không thể restart job chạy nhiều năm từ đầu).

Uses of Stream Processing

Truyền thống dùng cho monitoring: phát hiện gian lận thẻ tín dụng, trading system theo quy tắc giá, giám sát máy móc nhà máy, hệ thống quân sự/tình báo.

Complex event processing (CEP)

  • Từ thập niên 1990; giống regex nhưng tìm pattern các event trong stream. Mô tả bằng ngôn ngữ khai báo (SQL-like) hoặc GUI; engine duy trì state machine; khớp thì phát ra complex event.
  • Đảo vai trò query và dữ liệu: DB lưu dữ liệu lâu dài, query là tạm thời; CEP lưu query lâu dài, event chảy qua.
  • Ví dụ: Esper, IBM InfoSphere Streams, Apama, TIBCO StreamBase, SQLstream; Samza thêm SQL.

Stream analytics

  • Ít quan tâm chuỗi event cụ thể, thiên về aggregation và thống kê: tốc độ event theo thời gian, rolling average, so sánh với cùng kỳ tuần trước để phát hiện xu hướng/bất thường.
  • Tính trên window (ví dụ QPS trung bình và p99 trong 5 phút gần nhất).
  • Dùng probabilistic algorithm: Bloom filter, HyperLogLog (đếm cardinality), percentile estimation — tiết kiệm bộ nhớ. Nhưng stream processing không vốn dĩ là gần đúng; xấp xỉ chỉ là tối ưu.
  • Framework: Apache Storm, Spark Streaming, Flink, Concord, Samza, Kafka Streams; hosted: Google Cloud Dataflow, Azure Stream Analytics.

Maintaining materialized views

  • Giữ cache, search index, warehouse đồng bộ với DB nguồn = duy trì materialized view. State trong event sourcing cũng là materialized view.
  • Khác analytics: cần mọi event từ đầu đến giờ (window kéo về "đầu thời gian"), trừ event đã bị compaction bỏ. Samza, Kafka Streams hỗ trợ nhờ log compaction của Kafka.

Search on streams

  • Tìm từng event theo tiêu chí phức tạp (full-text): dịch vụ media monitoring theo dõi tin tức nhắc tới công ty/sản phẩm; website bất động sản báo khi có nhà khớp tiêu chí. Elasticsearch percolator.
  • Giống CEP: query được lưu, document chảy qua. Nhiều query thì có thể index cả query để thu hẹp tập query cần kiểm tra.

Message passing and RPC

Actor model/message passing không được coi là stream processing: - Actor chủ yếu quản lý concurrency và thực thi phân tán; stream processing là kỹ thuật quản lý dữ liệu. - Giao tiếp actor thường ngắn hạn, one-to-one; event log durable, multi-subscriber. - Actor giao tiếp tùy ý (kể cả vòng request/response); stream processor thường là pipeline acyclic.

Có vùng giao thoa: Storm distributed RPC (query xen kẽ với event stream rồi gom kết quả). Có thể xử lý stream bằng actor framework, nhưng nhiều framework không đảm bảo giao message khi crash.

Reasoning About Time

  • "Trung bình 5 phút gần nhất" nghe rõ ràng nhưng thực ra rất khó.
  • Batch: nhìn timestamp trong event, không nhìn đồng hồ máy chạy job (xử lý 1 năm dữ liệu trong vài phút) → deterministic.
  • Nhiều stream framework dùng đồng hồ local của máy xử lý (processing time) — đơn giản, ổn nếu độ trễ không đáng kể, nhưng hỏng khi có processing lag.

Event time versus processing time

  • Nguyên nhân trễ: queueing, lỗi mạng, contention ở broker, restart consumer, reprocess event cũ.
  • Trễ còn gây đảo thứ tự: user gửi request 1 (server A) rồi request 2 (server B); event của B tới broker trước A.
  • Ví dụ Star Wars: Episode IV (1977), V, VI rồi mới I, II, III (1999–2005), VII (2015). Số episode ≈ event time, ngày bạn xem ≈ processing time.
  • Nhầm lẫn gây dữ liệu sai: redeploy stream processor, dừng 1 phút rồi xử lý backlog → nếu đo tốc độ request theo processing time sẽ thấy spike giả (Figure 11-7), dù tốc độ thật vẫn đều.

Knowing when you're ready

  • Window theo event time: không bao giờ chắc đã nhận đủ event của window. Đang đếm phút 37, đa số event giờ thuộc phút 38–39 — khi nào đóng window 37?
  • Timeout rồi đóng, nhưng vẫn có straggler tới muộn (bị buffer ở máy khác vì mạng). Hai lựa chọn:
    1. Bỏ qua straggler; theo dõi số event bị bỏ như metric, alert nếu nhiều.
    2. Phát correction — giá trị cập nhật có straggler, có thể phải retract output cũ.
  • Có thể dùng message đặc biệt "từ giờ không còn message nào có timestamp < t" (ý tưởng watermark), nhưng nhiều producer thì phải theo dõi từng producer.

Whose clock are you using, anyway?

  • Mobile app dùng offline, buffer event, gửi sau vài giờ/ngày → straggler cực trễ.
  • Timestamp nên là lúc tương tác xảy ra theo đồng hồ thiết bị, nhưng đồng hồ thiết bị user không tin được; đồng hồ server chính xác hơn nhưng ít ý nghĩa.
  • Giải pháp: log 3 timestamp:
    1. Thời điểm event xảy ra (đồng hồ thiết bị).
    2. Thời điểm gửi lên server (đồng hồ thiết bị).
    3. Thời điểm server nhận (đồng hồ server).
  • (3) − (2) ≈ độ lệch đồng hồ thiết bị/server (bỏ qua network delay) → áp vào (1) để ước lượng thời điểm thật.
  • Batch cũng gặp vấn đề này, chỉ là stream làm ta nhận ra rõ hơn.

Types of windows

Loại windowĐộ dàiChồng lấn?Ví dụ / cách implement
TumblingCố địnhKhông — mỗi event thuộc đúng 1 window1 phút: 10:03:00–10:03:59, 10:04:00–10:04:59… Làm tròn timestamp xuống phút
HoppingCố định, có hop sizeCó — để làm mượt5 phút, hop 1 phút: 10:03–10:07:59, 10:04–10:08:59… Tính tumbling 1 phút rồi gộp các window kề nhau
SlidingKhoảng cách giữa các eventCó, không có biên cố định5 phút: event 10:03:39 và 10:08:12 cùng window vì cách nhau < 5 phút. Giữ buffer event theo thời gian, loại event hết hạn
SessionKhông cố địnhKhôngGom event của cùng user gần nhau; đóng khi user không hoạt động (ví dụ 30 phút). Dùng cho website analytics

Stream Joins

Stream processing tổng quát hóa data pipeline cho dữ liệu unbounded → cũng cần join, nhưng khó hơn vì event mới có thể đến bất cứ lúc nào. Ba loại:

Stream-stream join (window join)

  • Ví dụ click-through rate: mỗi lần search log một event (query + kết quả), mỗi lần click log một event; nối qua session ID. Tương tự trong quảng cáo.
  • Click có thể không bao giờ đến (bỏ search), hoặc đến sau vài giây tới vài tuần; thậm chí click tới trước search do network delay → chọn window, ví dụ join nếu cách nhau ≤ 1 giờ.
  • Nhúng thông tin search vào click event không tương đương join: chỉ biết các search được click, không biết search nào không được click → CTR sai.
  • Implement: giữ state — mọi event trong 1 giờ qua, index theo session ID. Event đến thì thêm vào index và kiểm tra index bên kia; khớp thì phát event "kết quả nào được click"; search hết hạn mà không có click thì phát event "không được click".

Stream-table join (stream enrichment)

  • Ví dụ: stream user activity event (có user ID) + DB user profile → output activity event được làm giàu với thông tin profile.
  • Query DB từ xa cho từng event: chậm, có thể làm quá tải DB.
  • Tốt hơn: nạp bản sao DB vào stream processor (hash table in-memory hoặc index trên disk local) — giống map-side hash join trong batch.
  • Khác batch: batch dùng snapshot tại một thời điểm, stream processor chạy lâu dài nên bản sao phải được cập nhật → subscribe changelog của DB profile qua CDC. Thực chất là join giữa hai stream: activity event và profile update.
  • So với stream-stream: phía changelog dùng window "từ đầu thời gian" (vô hạn), bản mới ghi đè bản cũ; phía activity có thể không cần window.

Table-table join (materialized view maintenance)

  • Ví dụ Twitter timeline: duy trì "inbox" timeline cache cho mỗi user để đọc chỉ là một lookup:
    • User u tweet → thêm vào timeline mọi follower của u.
    • Xóa tweet → xóa khỏi mọi timeline.
    • u1 follow u2 → thêm tweet gần đây của u2 vào timeline u1.
    • u1 unfollow u2 → gỡ tweet của u2 khỏi timeline u1.
  • Cần stream tweet (gửi/xóa) và stream follow (follow/unfollow); processor giữ DB tập follower của mỗi user.
  • Đây là duy trì materialized view cho query join 2 bảng tweets JOIN follows ON followee_id = sender_id GROUP BY follower_id. Timeline là cache của kết quả query đó.
  • Chú thích thú vị: coi stream là đạo hàm của table, join là tích u·v → thay đổi của join theo quy tắc tích (u·v)' = u'v + uv': thay đổi tweets join với followers hiện tại, thay đổi followers join với tweets hiện tại.
Loại joinInputState cần giữVí dụOutput
Stream-stream (window join)Hai stream activity eventEvent trong window gần đây, index theo keySearch + click theo session ID (CTR)Stream các cặp khớp / không khớp
Stream-table (enrichment)Activity stream + changelog DBBản sao local của bảng (window vô hạn, cập nhật qua CDC)Activity event + user profileActivity event đã enrich
Table-table (materialized view)Hai changelog DBTrạng thái mới nhất của cả hai bảngTweets × follows → Twitter timelineStream thay đổi của materialized view

Time-dependence of joins

  • Cả 3 loại đều: giữ state từ một input, query state đó khi có message từ input kia. Thứ tự event cập nhật state rất quan trọng (follow rồi unfollow khác unfollow rồi follow).
  • Trong partition có thứ tự, giữa các stream/partition thì không → nếu user cập nhật profile, activity event nào join với profile cũ, event nào với profile mới?
  • Ví dụ thuế suất: hóa đơn phải join với thuế suất tại thời điểm bán, không phải hiện tại (quan trọng khi reprocess dữ liệu lịch sử).
  • Thứ tự không xác định → join nondeterministic: chạy lại cùng input có thể ra kết quả khác.
  • Trong data warehouse gọi là slowly changing dimension (SCD): gán ID riêng cho từng phiên bản record (mỗi lần đổi thuế suất có ID mới, hóa đơn lưu ID thuế suất lúc bán) → join deterministic, nhưng không log compaction được vì phải giữ mọi phiên bản.

Fault Tolerance

  • Batch chịu lỗi dễ: task lỗi thì chạy lại, output của task lỗi bị bỏ (input immutable, output chỉ hiện ra khi task xong). Kết quả như thể mỗi record được xử lý đúng một lần → exactly-once semantics (gọi là effectively-once thì chính xác hơn).
  • Stream không thể "đợi xong mới công bố output" vì stream vô hạn.

Microbatching and checkpointing

  • Microbatching (Spark Streaming): chia stream thành block nhỏ (~1 giây), mỗi block là một mini batch. Batch nhỏ → overhead scheduling lớn; batch lớn → trễ cao. Ngầm tạo tumbling window theo processing time bằng batch size; window lớn hơn phải tự mang state qua các batch.
  • Checkpointing (Apache Flink): định kỳ tạo rolling checkpoint của state ghi ra durable storage; operator crash thì khởi động lại từ checkpoint gần nhất, bỏ output sinh ra sau checkpoint. Checkpoint được kích hoạt bởi barrier trong stream, không ép kích thước window.
  • Hai cách cho exactly-once trong phạm vi framework. Nhưng khi output rời khỏi framework (ghi DB, gửi message ra broker ngoài, gửi email) thì không thể "bỏ" output → restart gây side effect lặp lại.

Atomic commit revisited

  • Cần mọi output và side effect của việc xử lý một event xảy ra khi và chỉ khi xử lý thành công: message gửi downstream, email/push, ghi DB, thay đổi operator state, ack input (bao gồm dịch offset). Tất cả atomic.
  • XA có nhiều vấn đề, nhưng trong môi trường hạn chế có thể làm hiệu quả: Google Cloud Dataflow, VoltDB, Kafka (transactions). Khác XA: không xuyên các công nghệ khác nhau, giữ transaction nội bộ trong framework; amortize overhead bằng cách gom nhiều message vào một transaction.

Idempotence

  • Idempotent: làm nhiều lần có hiệu quả như làm một lần. Set key = giá trị cố định là idempotent; tăng counter thì không.
  • Làm cho idempotent bằng metadata: ghi kèm offset Kafka của message gây ra lần ghi cuối vào DB ngoài → biết update đã được áp chưa. Storm Trident dùng ý tưởng tương tự.
  • Giả định cần có: restart phải replay cùng message theo cùng thứ tự (log-based broker làm được), xử lý deterministic, và không node nào khác cập nhật cùng giá trị đồng thời. Failover có thể cần fencing để chặn node "tưởng chết mà vẫn sống".
  • Dù nhiều điều kiện, đây là cách đạt exactly-once với overhead nhỏ.

Rebuilding state after a failure

Stateful operator (window aggregation, counter, bảng/index cho join) phải khôi phục được state: - Giữ state ở remote datastore có replication — nhưng query remote cho mỗi message thì chậm. - Giữ state local và replicate định kỳ: - Flink: snapshot operator state định kỳ ra HDFS. - Samza, Kafka Streams: gửi thay đổi state vào Kafka topic có log compaction (giống CDC). - VoltDB: xử lý mỗi message dư thừa trên nhiều node. - Đôi khi không cần replicate: window ngắn thì replay input là đủ nhanh; bản sao DB qua CDC thì dựng lại từ topic compacted. - Trade-off tùy hạ tầng: có hệ thống network nhanh hơn disk — không có lựa chọn tối ưu cho mọi tình huống.

Cơ chếCách hoạt độngVí dụƯuNhược
MicrobatchingChia stream thành batch ~1s, mỗi batch xử lý như batch jobSpark StreamingĐơn giản, tái dùng mô hình batchLatency ≥ batch size; window gắn với processing time; kém với hopping/sliding window
CheckpointingSnapshot state định kỳ qua barrier, restart từ checkpointApache FlinkLatency thấp, không ép windowKhông che được side effect ra ngoài framework
Atomic commit nội bộOutput, state, offset commit cùng một transactionGoogle Cloud Dataflow, VoltDB, Kafka transactionsExactly-once cả output lẫn stateChỉ trong phạm vi một hệ thống; có overhead
IdempotenceGhi kèm offset/ID để phát hiện ghi trùngStorm Trident, ghi DB kèm offset KafkaOverhead nhỏ, dùng được với hệ thống ngoàiCần replay đúng thứ tự, xử lý deterministic, fencing khi failover

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

  • "Kafka là message queue như RabbitMQ" — sai về bản chất: Kafka là log bền vững, đọc không xóa, replay được, load balancing theo partition. Chọn sai loại broker dẫn tới thiết kế sai (ví dụ cần song song theo từng message mà lại dùng log với ít partition).
  • Dual writes "trông có vẻ ổn" — không có lỗi nào vẫn có thể lệch dữ liệu vĩnh viễn do race condition; cộng thêm partial failure. Dùng CDC / một log làm nguồn thứ tự duy nhất.
  • Giả định Kafka giữ thứ tự toàn cục — chỉ có thứ tự trong partition. Muốn thứ tự cho một entity thì partition theo key của entity đó.
  • Tăng consumer quá số partition để tăng tốc — consumer thừa sẽ ngồi không.
  • Dùng processing time cho analytics — gây spike giả khi restart/backlog, sai khi reprocess. Dùng event time + xử lý straggler.
  • Tin đồng hồ thiết bị người dùng — có thể sai (cố ý hoặc vô tình); dùng kỹ thuật 3 timestamp.
  • Nghĩ "exactly-once" nghĩa là message chỉ được xử lý đúng một lần — thực tế là effectively-once: có thể xử lý nhiều lần nhưng hiệu ứng quan sát được như một lần. Và nó thường chỉ đúng bên trong framework; side effect ra ngoài (email, API call) cần idempotence hoặc atomic commit.
  • Nhúng dữ liệu search vào click event thay cho join — mất thông tin về search không được click.
  • Join không quan tâm tới thời gian — join với state "hiện tại" khi reprocess dữ liệu cũ cho kết quả sai (thuế suất, tỷ giá); cần versioned record (SCD).
  • Consumer của event log reject event — không được; validation phải trước khi command thành event.
  • "Immutable nên không bao giờ xóa được" — pháp luật (GDPR) có thể bắt buộc xóa; cần thiết kế excision/crypto-shredding từ đầu.
  • Stream processing luôn là xấp xỉ — sai; thuật toán xác suất chỉ là tối ưu tùy chọn.
  • Event sourcing = CDC — khác mức trừu tượng; event sourcing không log compact được theo key.

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

  • Đồng bộ DB với Elasticsearch/Redis/warehouse: câu trả lời chuẩn là CDC (Debezium) → Kafka → consumer ghi vào index/cache, không dual write. Nhắc: log compaction để bootstrap consumer mới, initial snapshot gắn offset, replication lag → cân nhắc read-your-writes.
  • Transactional outbox pattern (thực tế rất phổ biến): ghi business data và event vào bảng outbox trong cùng một transaction DB, rồi CDC đọc outbox đẩy lên Kafka — giải quyết dual-write mà không cần 2PC.
  • Chọn broker trong phỏng vấn: task queue (gửi email, resize ảnh, job tốn công, không cần thứ tự) → RabbitMQ/SQS; event stream, activity tracking, CDC, pipeline analytics, cần replay → Kafka/Kinesis. Giải thích partition key quyết định thứ tự và mức song song.
  • Thiết kế real-time analytics (đếm view, top-K, trending hashtag, ad click aggregation): Kafka + Flink/Kafka Streams; tumbling/hopping window theo event time; watermark + allowed lateness cho straggler; HyperLogLog cho unique visitors; idempotent sink hoặc exactly-once để không đếm trùng.
  • Ad click / CTR: ví dụ stream-stream join kinh điển trong phỏng vấn "Design ad click aggregator".
  • News feed / Twitter timeline: fan-out-on-write chính là table-table join duy trì materialized view; nhắc hybrid cho celebrity.
  • Enrichment (gắn profile, tỷ giá, geo vào event): stream-table join với bản sao local cập nhật qua CDC thay vì gọi RPC cho mỗi event.
  • Fraud detection, monitoring, alerting: CEP / pattern trên stream.
  • Event sourcing + CQRS: ngân hàng, ledger, đặt vé, giỏ hàng — nhấn mạnh audit, replay, dựng view mới; đánh đổi eventual consistency và độ phức tạp.
  • Vận hành: monitor consumer lag, retention theo ngày/tuần, số partition quyết định giới hạn song song, rebalancing khi consumer chết gây xử lý lại → consumer phải idempotent.

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

Bấm vào câu hỏi để xem đáp án
1Khác biệt cốt lõi giữa AMQP/JMS-style broker và log-based broker là gì? Khi nào chọn cái nào?
AMQP/JMS gán từng message cho consumer, xóa khi ack, không replay, có thể đảo thứ tự khi redelivery — hợp với task queue tốn công, không cần thứ tự. Log-based append vào log partition, consumer đọc tuần tự theo offset, message được giữ lại nên replay được, có thứ tự trong partition, song song theo partition — hợp với throughput cao, cần thứ tự, derived data.
2Vì sao dual writes gây inconsistency ngay cả khi không có lỗi nào?
Hai client ghi đồng thời, DB và search index có thể nhận hai write theo thứ tự khác nhau (race condition, Figure 11-4) → giá trị cuối khác nhau vĩnh viễn. Ngoài ra một write có thể thành công, write kia thất bại. Gốc rễ là không có một leader duy nhất quyết định thứ tự.
3Log compaction giúp gì cho CDC, và tại sao không dùng được như vậy với event sourcing?
Với CDC mỗi event chứa toàn bộ giá trị mới của key, nên chỉ giữ event mới nhất mỗi key là đủ để dựng lại DB từ offset 0 mà không cần snapshot. Event sourcing lưu ý định (hành động), event sau không ghi đè event trước, nên cần toàn bộ lịch sử — dùng snapshot thay thế.
4Phân biệt command và event trong event sourcing. Tại sao consumer không được reject event?
Command là request có thể thất bại, cần validate; khi được chấp nhận nó trở thành event — một fact immutable, bền vững. Consumer thấy event khi nó đã nằm trong log và có thể đã được consumer khác xử lý, nên validation phải xảy ra đồng bộ trước khi ghi event (hoặc tách thành tentative + confirmation).
5Vì sao dùng processing time cho windowing có thể cho kết quả sai? Xử lý straggler thế nào?
Khi có lag (restart, backlog, reprocess), event bị xử lý dồn lại → tạo spike giả và gán sai window. Dùng event time; straggler thì hoặc bỏ qua (kèm metric), hoặc phát correction/retract; có thể dùng watermark "không còn event < t".
6Nêu 4 loại window và khác biệt chính.
Tumbling — cố định, không chồng lấn; hopping — cố định, chồng lấn theo hop size; sliding — gom event cách nhau trong khoảng cho trước, không biên cố định; session — không độ dài cố định, đóng sau khoảng không hoạt động của user.
7Mô tả 3 loại stream join kèm ví dụ, và "time-dependence of joins" là gì?
Stream-stream (search + click trong 1 giờ để tính CTR), stream-table (activity event enrich bằng user profile với bản sao local cập nhật qua CDC), table-table (tweets × follows duy trì Twitter timeline). Time-dependence: join với state thay đổi theo thời gian, thứ tự giữa các stream không xác định nên kết quả có thể nondeterministic; giải bằng versioned ID (slowly changing dimension), đổi lại không log compact được.
8Các cách đạt exactly-once trong stream processing và giới hạn của chúng?
Microbatching (Spark), checkpointing (Flink) — chỉ đúng trong framework; atomic commit nội bộ (Dataflow, VoltDB, Kafka transactions) — gộp output, state, offset vào một transaction; idempotence (ghi kèm offset) — cần replay đúng thứ tự, deterministic, fencing. Side effect ra ngoài framework cần idempotence hoặc atomic commit.

Đâ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.