Reactive Streams và điều khiển áp lực ngược (Backpressure)
1. Tổng quan
A. Định nghĩa
Reactive Streams là đặc tả và mô hình lập trình chuẩn hóa tín hiệu áp lực ngược (Backpressure), nhằm xử lý luồng dữ liệu bất đồng bộ·không chặn (non-blocking) giữa bên sản xuất và bên tiêu thụ, nhưng chỉ cho dữ liệu chảy với tốc độ mà bên tiêu thụ có thể đáp ứng.
Bản chất của Reactive Streams nằm ở việc kết hợp "đẩy (push)" và "kéo (pull)" vào một giao thức duy nhất.
Xử lý hướng sự kiện truyền thống là mô hình push trong đó bên sản xuất đơn phương đẩy dữ liệu sang bên tiêu thụ, nên nếu bên tiêu thụ chậm, bộ đệm sẽ tràn hoặc thông điệp bị mất.
Ngược lại, mô hình pull — bên tiêu thụ yêu cầu dữ liệu mỗi khi cần — thì an toàn nhưng độ trễ khứ hồi lớn khiến thông lượng giảm.
Reactive Streams đưa vào tín hiệu request(n) để bên tiêu thụ thông báo nhu cầu (demand) "hiện có thể nhận tối đa n phần tử" cho bên sản xuất, qua đó hiện thực mô hình push-pull động: bình thường giữ hiệu quả của push, còn khi bên tiêu thụ bão hòa thì luồng tự động được điều khiển.
Ở đây, áp lực ngược là tín hiệu điều khiển trả phần thiếu hụt ngược về thượng nguồn để làm chậm chính tốc độ sản xuất khi năng lực xử lý của bên tiêu thụ thấp hơn tốc độ sản xuất của bên sản xuất. Giống như khi hạ lưu ống nước bị tắc thì áp lực nước ở thượng lưu tăng lên, trong pipeline phần mềm, độ trễ ở hạ nguồn cũng phải được truyền lên thượng nguồn thì toàn hệ thống mới không sụp đổ. Hệ thống không có áp lực ngược, khi tải dồn về, sẽ rơi vào một trong hai kiểu thất bại: hàng đợi phình vô hạn cho tới khi tiến trình bị kết thúc vì thiếu bộ nhớ, hoặc nếu hàng đợi hữu hạn thì thông điệp bị âm thầm bỏ đi.
B. Bối cảnh ra đời và sự cần thiết
Giai đoạn 2013~2015 khi Reactive Streams được xác lập thành chuẩn là thời kỳ dữ liệu streaming và microservice bùng nổ.
Nhiều thư viện như RxJava, Akka Streams, Project Reactor xử lý luồng bất đồng bộ theo cách riêng, và khi kết hợp các thư viện khác nhau thì vấn đề tín hiệu áp lực ngược không tương thích lặp đi lặp lại.
Để giải quyết, Netflix, Lightbend, Pivotal, v.v. đã cùng tham gia thống nhất một đặc tả tối thiểu gồm bốn interface Publisher, Subscriber, Subscription, Processor, sau đó được đưa vào thư viện chuẩn dưới dạng java.util.concurrent.Flow của Java 9.
Sự cần thiết được giải thích theo hai trục: hiệu quả tài nguyên và tính ổn định. Thứ nhất, mô hình mỗi-yêu-cầu-một-luồng (thread-per-request) chiếm một luồng cho mỗi yêu cầu, nên khi số kết nối đồng thời lên tới hàng vạn, bộ nhớ stack của luồng và chi phí chuyển ngữ cảnh tăng vọt. Ví dụ, nếu một luồng dùng stack 1MB thì 1 vạn luồng tiêu tốn 10GB, và vì phần lớn luồng đang chờ I/O nên CPU nhàn rỗi trong khi bộ nhớ cạn kiệt. Mô hình reactive không chặn xử lý hàng vạn kết nối đồng thời bằng một số ít luồng event loop (thường bằng số lõi CPU), tiết kiệm đáng kể bộ nhớ.
Thứ hai, trong pipeline phân tán, nếu độ trễ của một đoạn không được truyền lên thượng nguồn thì sự cố sẽ khuếch đại. Nếu Kafka consumer đã chậm lại do trễ ghi cơ sở dữ liệu mà vẫn tiếp tục kéo dữ liệu từ broker, heap của consumer sẽ tràn. Áp lực ngược khiến consumer trong tình huống này chỉ poll lượng có thể xử lý, để độ trễ được toàn bộ pipeline chia nhau hấp thụ. Do đó, Reactive Streams không chỉ là tối ưu hiệu năng đơn thuần mà là phương tiện cốt lõi của thiết kế khả năng phục hồi (resilience) trong tình huống tải tăng đột biến.
C. Đặc điểm
Đặc điểm của Reactive Streams được tóm gọn là tính bất đồng bộ, không chặn, điều khiển luồng dựa trên áp lực ngược và khả năng kết hợp (composability). Pipeline được mô tả theo kiểu khai báo bằng cách nối các toán tử (operator) theo phong cách hàm, và mỗi toán tử kiêm luôn vai trò điều hòa nhu cầu của thượng nguồn và hạ nguồn. Bảng dưới đây tóm tắt khác biệt với các mô hình truyền thống, còn "tại sao" của từng mục sẽ được giải thích bằng văn xuôi ở các mục tiếp theo.
| Phân loại | Mô hình đồng bộ chặn | Mô hình sự kiện push thuần | Reactive Streams |
|---|---|---|---|
| Điều khiển luồng | Call stack là áp lực ngược tự nhiên | Không có (nguy cơ mất·OOM) | Áp lực ngược tường minh dựa trên request(n) |
| Sử dụng luồng | Chiếm luồng theo yêu cầu | Event loop | Event loop + không chặn |
| Lan truyền độ trễ | Lan truyền qua lời gọi đồng bộ | Không lan truyền | Lan truyền qua tín hiệu nhu cầu |
| Khả năng kết hợp | Thấp | Trung bình | Cao (chuỗi toán tử) |
2. Đặc tả Reactive Streams và các thành phần
A. Bốn interface cốt lõi
Đặc tả Reactive Streams định nghĩa quy ước tương tác bằng bốn interface Publisher, Subscriber, Subscription, Processor.
Publisher là nguồn sản xuất dữ liệu, kết nối bên tiêu thụ bằng lời gọi subscribe().
Subscriber tiêu thụ dữ liệu và có bốn callback onSubscribe, onNext, onError, onComplete.
Subscription biểu thị một lần kết nối giữa bên sản xuất và bên tiêu thụ, cung cấp request(n) để bên tiêu thụ yêu cầu nhu cầu và cancel() để ngắt kết nối.
Processor là bước trung gian vừa là Publisher vừa là Subscriber, tương ứng với mỗi node trong chuỗi toán tử.
Biểu diễn cấu trúc tổng thể bằng sơ đồ khái niệm như sau.
flowchart LR
P["Publisher(bên sản xuất)"] -->|onSubscribe| SUB["Subscription(đăng ký)"]
SUB -->|request n| P
P -->|onNext dữ liệu| PR["Processor(chuỗi toán tử)"]
PR -->|onNext biến đổi| S["Subscriber(bên tiêu thụ)"]
S -->|request n nhu cầu| PR
PR -->|request n điều chỉnh lại| P
S -.->|onError / onComplete| END["Tín hiệu kết thúc"]
Cốt lõi của cấu trúc này là trong khi dữ liệu chảy từ trái sang phải, tín hiệu nhu cầu (request(n)) lại chảy ngược từ phải sang trái.
Nếu bên tiêu thụ chỉ yêu cầu số lượng mình đáp ứng được, tín hiệu đó đi qua các toán tử tới bên sản xuất, và bên sản xuất chỉ phát onNext không vượt quá số được yêu cầu.
Nhờ quy ước này, dù không biết năng lực xử lý của bên tiêu thụ, bên sản xuất tuyệt đối không đẩy dữ liệu vượt lượng yêu cầu.
B. Quy ước tín hiệu và bất biến
Đặc tả đòi hỏi các bất biến nghiêm ngặt về thứ tự và số lượng tín hiệu.
onSubscribe phải được gọi trước đúng một lần, onNext không được vượt quá số được yêu cầu, và sau onError hoặc onComplete không được phát sinh bất kỳ tín hiệu nào.
Các quy ước này là hợp đồng tối thiểu để kết nối an toàn các thư viện khác nhau, và việc tuân thủ được kiểm chứng bằng bộ kiểm thử chuẩn gọi là TCK (Technology Compatibility Kit).
Hiện thực không tuân thủ quy ước sẽ gây điều kiện tranh đua hoặc rò rỉ tài nguyên khi kết hợp, nên khuyến nghị thực tiễn là dùng các factory method của thư viện đã được kiểm chứng thay vì tự hiện thực Publisher.
C. Cold stream và hot stream
Luồng được chia thành cold và hot tùy theo thời điểm đăng ký. Cold stream bắt đầu sản xuất dữ liệu lại từ đầu mỗi khi có subscriber gắn vào, phù hợp với các công việc hoàn tất theo đơn vị yêu cầu như yêu cầu HTTP, đọc tệp, truy vấn cơ sở dữ liệu. Hot stream thì dữ liệu chảy liên tục bất kể có subscriber hay không, tương ứng với các nguồn mang tính phát sóng thời gian thực như sự kiện cảm biến, giá chứng khoán, cú nhấp của người dùng. Hot stream khó áp dụng nguyên áp lực ngược, vì nguồn sản xuất không chờ nhu cầu của bên tiêu thụ. Trong trường hợp này, phải thiết kế kèm các chiến lược cho phép mất dữ liệu như đệm, lấy mẫu, giữ giá trị mới nhất sẽ được giải thích ở phần sau.
3. Cơ chế điều khiển áp lực ngược và quy trình hoạt động
A. Quy trình điều khiển luồng dựa trên nhu cầu
Cơ chế chuẩn của áp lực ngược là cuộc đối thoại khứ hồi trong đó bên tiêu thụ gọi request(n) theo dư địa xử lý, và bên sản xuất chỉ phát trong phạm vi đó.
Bên tiêu thụ có thể chọn yêu cầu một nhu cầu lớn một lần rồi bổ sung mỗi khi tiêu hao (ví dụ: yêu cầu 256 phần tử, khi tiêu hao một nửa thì thêm 128), hoặc yêu cầu đúng từng phần tử một để khớp với tốc độ xử lý.
Sơ đồ tuần tự dưới đây cho thấy nhu cầu được điều chỉnh thế nào khi bên tiêu thụ chậm lại.
sequenceDiagram
participant S as Subscriber(bên tiêu thụ)
participant O as Operator(bộ đệm)
participant P as Publisher(bên sản xuất)
S->>O: request(256)
O->>P: request(256)
P-->>O: onNext x256
O-->>S: onNext x256
Note over S: Phát sinh trễ xử lý
S->>O: chỉ bổ sung request(64)
O->>P: request chỉ bằng phần trống của bộ đệm
Note over P: Tốc độ phát tự động giảm
S->>O: cancel() hoặc hoàn tất
Điểm đáng chú ý trong quy trình này là khi bên tiêu thụ bị trễ, yêu cầu bổ sung đến muộn, và kết quả là việc phát của bên sản xuất cũng tự nhiên chậm lại.
Không cần bất kỳ "lệnh giới hạn tốc độ" tường minh nào, chỉ riêng độ trễ của tín hiệu nhu cầu đã khiến tốc độ toàn pipeline khớp với bên tiêu thụ.
Trong thực tế, Flux của Project Reactor mặc định đặt bộ đệm prefetch kích thước 256 và yêu cầu lô tiếp theo khi 75% đã tiêu hao, tự động hóa cuộc đối thoại này.
B. Kịch bản sụp đổ khi không có áp lực ngược
Nếu áp lực ngược không được truyền đúng, hệ thống sụp đổ theo một trong hai hướng.
Dùng bộ đệm vô hạn thì khi tải tăng đột biến, hàng đợi tiếp tục phình ra làm cạn heap và tiến trình kết thúc với OutOfMemoryError.
Dùng bộ đệm hữu hạn và ném ngoại lệ khi tràn thì yêu cầu thất bại với lỗi như MissingBackpressureException.
Trong các phân tích hậu sự cố của nhiều dịch vụ streaming đầu những năm 2020, các trường hợp độ trễ tiêu thụ không truyền được lên thượng nguồn khiến hàng đợi giữa broker–consumer bùng nổ đã được báo cáo lặp lại.
Do đó, thiết kế áp lực ngược phải được kiểm chứng không theo tiêu chí "có chạy tốt ở tải bình thường không" mà theo "khi tải tối đa vượt năng lực tiêu thụ thì suy giảm một cách nhẹ nhàng (graceful) ra sao".
C. Chiến lược cho phép mất dữ liệu và đệm
Với hot stream không thể chờ bên tiêu thụ vô hạn, cần chiến lược tràn (overflow) cố ý bỏ hoặc tóm tắt một phần dữ liệu. Tiêu biểu, buffer tích lũy tới một kích thước nhất định rồi chuyển theo lô, drop bỏ dữ liệu mới khi bên tiêu thụ bận, còn latest chỉ giữ giá trị gần nhất và ghi đè giá trị trước. Lấy mẫu (sample) hay cửa sổ (window) gom dữ liệu theo tiêu chí thời gian·số lượng để giảm tải hạ nguồn. Bảng dưới đây so sánh đặc tính của từng chiến lược.
| Chiến lược | Hoạt động | Mất dữ liệu | Tình huống phù hợp |
|---|---|---|---|
| buffer | Nạp vào hàng đợi hữu hạn rồi chuyển theo lô | Tràn thì lỗi/chặn | Tải nhất thời ngắn, không được mất |
| drop | Bỏ phần vượt | Có | Khi thứ tự quan trọng hơn tính mới |
| latest | Chỉ giữ 1 bản ghi mới nhất | Có | Dashboard·hiển thị giá |
| error | Phát tín hiệu thất bại ngay | - | Batch mà thất bại nhanh là tốt hơn |
Việc chọn chiến lược xuất phát từ yêu cầu nghiệp vụ. Dữ liệu không được mất dù một bản ghi như sự kiện thanh toán phải bảo đảm không tổn thất bằng buffer + hàng đợi bền + xử lý lại, còn dữ liệu chỉ giá trị mới nhất có ý nghĩa như đồng hồ nhiệt độ thời gian thực thì dùng latest để bỏ giá trị cũ sẽ có lợi cả về hiệu quả tài nguyên lẫn trải nghiệm người dùng.
4. So sánh và tình huống áp dụng thực tế
A. Đánh đổi so với mô hình mệnh lệnh
Mô hình reactive mang lại thông lượng cao và hiệu quả tài nguyên nhưng phải trả giá bằng khó gỡ lỗi và đường cong học tập. Trong pipeline bất đồng bộ, stack trace không chứa đường thực thi thực tế nên khó truy nguyên nhân lỗi, và chỉ một lời gọi chặn làm tắc event loop là toàn bộ thông lượng tụt mạnh. Ngược lại, mô hình thread-per-request có mã trực quan và dễ gỡ lỗi, và gần đây luồng ảo (Virtual Thread, Project Loom) của Java xuất hiện giúp đạt độ đồng thời cao ngay cả với mã chặn. Do đó, điều quan trọng là phán đoán áp dụng có chọn lọc cho các đoạn vừa bị giới hạn bởi I/O vừa về bản chất cần streaming·điều khiển áp lực ngược, chứ không phải "reactive bằng mọi giá".
B. Tình huống dịch vụ web — Spring WebFlux
Giả sử một API gateway phải xử lý hàng vạn yêu cầu mỗi giây bằng một số ít event loop.
Với Spring MVC (thread-per-request), 200 luồng trong thread pool của Tomcat đều bị chiếm khi chờ phản hồi của API bên ngoài, dẫn tới hiện tượng dù mức dùng CPU thực tế chỉ 20% mà yêu cầu vẫn chờ trong hàng đợi rồi hết thời gian chờ.
Khi chuyển sang Spring WebFlux (dựa trên Reactor Netty), số event loop bằng số lõi CPU xử lý không chặn hàng vạn kết nối, và khi API hạ nguồn chậm lại thì request(n) được truyền lên client thượng nguồn, tự nhiên làm chậm luồng vào.
Trong các ca đo thực tế, người ta báo cáo khả năng tiếp nhận kết nối đồng thời tăng vài lần trên cùng phần cứng và độ trễ P99 ổn định hơn.
C. Tình huống pipeline dữ liệu — Kafka consumer
Giả sử trong pipeline sự kiện dựa trên Kafka, consumer có thể ghi 5.000 bản ghi mỗi giây vào cơ sở dữ liệu hạ nguồn nhưng topic nhận vào 20.000 bản ghi mỗi giây.
Nếu poll không giới hạn mà không có áp lực ngược, các bản ghi chưa xử lý chất đống trong heap của consumer dẫn tới bùng nổ GC và OOM.
Connector Kafka của Reactor Kafka hay Akka Streams điều chỉnh tốc độ gọi poll và commit offset theo nhu cầu của sink hạ nguồn, đồng thời giới hạn prefetch theo từng phân vùng, để consumer chỉ kéo lượng có thể xử lý.
Kết quả là dù luồng vào tức thời vượt năng lực tiêu thụ, dữ liệu vẫn được lưu an toàn tại broker trong khi consumer bắt kịp với tốc độ ổn định, và toàn pipeline đệm được đợt dồn dập đầu vào.
5. Chuyên sâu — Xu hướng mới và liên kết với công nghệ tương tự
Chuẩn Reactive Streams đã được đưa vào API Flow của Java 9 và trở thành một phần của chuẩn ngôn ngữ, đồng thời áp lực ngược đang mở rộng ra mọi tầng của stack như R2DBC (truy cập cơ sở dữ liệu quan hệ theo kiểu reactive) và reactive gRPC.
Nếu driver cơ sở dữ liệu không hỗ trợ áp lực ngược, một điểm trong pipeline sẽ bị chặn và toàn bộ lợi ích sụp đổ, nên "không chặn đầu-cuối (end-to-end)" đã trở thành mối quan tâm cốt lõi trong thực tế.
Mặt khác, luồng ảo được chính thức hóa trong Java 21 đã nổi lên như một phương án thay thế cho mô hình reactive. Luồng ảo chạy mã trông như chặn theo kiểu không chặn, cho phép đạt hiệu quả tài nguyên của reactive bằng mã mệnh lệnh, nhưng điều quan trọng là bản thân luồng ảo không cung cấp áp lực ngược. Nghĩa là vấn đề đồng thời (thiếu luồng) thì luồng ảo giải quyết, còn vấn đề điều khiển luồng (bất tương xứng tốc độ sản xuất-tiêu thụ) vẫn cần cơ chế riêng như Reactive Streams hoặc semaphore, hàng đợi giới hạn. Vì vậy, xu hướng là hai công nghệ không cạnh tranh mà định hình thành quan hệ bổ trợ: "luồng ảo cho I/O đơn giản, Reactive Streams cho streaming cần áp lực ngược".
Khái niệm áp lực ngược không chỉ giới hạn trong Reactive Streams. Cửa sổ nhận của TCP, điều khiển luồng HTTP/2 của gRPC, điều tiết dựa trên consumer lag của Kafka, việc chặn tải của circuit breaker đều thuộc họ áp lực ngược·điều khiển luồng theo nghĩa rộng. Từ góc nhìn Kỹ sư chuyên nghiệp, nếu gom chúng vào một phả hệ "điều khiển luồng" để giải thích, có thể xây dựng bài làm xoay quanh nguyên lý thay vì học thuộc từng công nghệ riêng lẻ.
6. Các điểm cần cân nhắc và hàm ý
Thứ nhất, phải phán đoán tỉnh táo phạm vi áp dụng. Reactive Streams mang lại lợi ích lớn ở các đoạn giới hạn bởi I/O, độ đồng thời cao và streaming, nhưng với batch giới hạn bởi CPU hay CRUD đơn giản thì chỉ làm tăng độ phức tạp. Khi luồng ảo đã trưởng thành như hiện nay, cần tiêu chí ra quyết định: trước hết phân biệt "vấn đề là mã chặn hay là điều khiển luồng", và chỉ chọn reactive khi là trường hợp sau.
Thứ hai, phải lấy không chặn đầu-cuối làm nguyên tắc thiết kế. Chỉ cần ở một điểm bất kỳ trong pipeline có lời gọi chặn (JDBC, I/O tệp đồng bộ, v.v.) chạy trên event loop, số ít event loop sẽ bị tắc và toàn bộ thông lượng sụp đổ. Các công việc chặn không tránh được nhất định phải được cô lập vào thread pool biên riêng (bounded elastic scheduler), và tầng truy cập dữ liệu nên được thống nhất bằng driver không chặn như R2DBC.
Thứ ba, phải chọn tường minh chiến lược tràn phù hợp với mức rủi ro nghiệp vụ. Miền không cho phép mất dữ liệu thì kết hợp bộ đệm không tổn thất với hàng đợi bền·xử lý lại, miền chỉ cần tính mới thì dùng chiến lược cho phép mất để tiết kiệm tài nguyên; và bản thân chính sách tràn nên được lưu thành hồ sơ quyết định kiến trúc (ADR) để người vận hành có thể hiểu.
Thứ tư, phải thiết kế đồng thời khả năng quan sát và kiểm thử. Áp lực ngược không hiện ra ở tải bình thường mà chỉ lộ ra ở tải tối đa, nên cần thu thập thường xuyên các chỉ số như kích thước prefetch, tỷ lệ chiếm dụng bộ đệm, lượng phát so với yêu cầu, consumer lag, đồng thời xác nhận trước "đường cong suy giảm khi vượt năng lực tiêu thụ" bằng kiểm thử tải và tiêm lỗi. Checkpoint và lan truyền ngữ cảnh để bổ sung cho stack trace bất đồng bộ cũng là một phần của chuẩn bị vận hành.
Thứ năm, phải song hành năng lực tổ chức và chuẩn hóa.
Mã reactive có đường cong học tập dốc và nếu dùng sai còn làm hiệu năng tệ hơn, nên an toàn hơn cả là quy định các thành ngữ (idiom) của thư viện đã kiểm chứng làm chuẩn của nhóm, và đưa vào code review quy tắc kết hợp factory·toán tử thay vì tự hiện thực Publisher.
Tài liệu tham khảo
- Reactive Streams Specification: https://www.reactive-streams.org/
- Project Reactor Reference — Handling Backpressure: https://projectreactor.io/docs/core/release/reference/
- Spring Framework, Web on Reactive Stack (WebFlux): https://docs.spring.io/spring-framework/reference/web/webflux.html
- OpenJDK, JEP 444: Virtual Threads: https://openjdk.org/jeps/444
Tóm tắt một câu: Reactive Streams là chuẩn xử lý luồng bất đồng bộ, dùng tín hiệu áp lực ngược dựa trên
request(n)để điều khiển ngược tốc độ sản xuất theo năng lực xử lý của bên tiêu thụ, qua đó đồng thời đạt hiệu quả tài nguyên của mô hình không chặn và khả năng phục hồi khi tải tăng đột biến.