RxJava: kiến thức cơ bản, ReactiveX và làm việc với luồng dữ liệu

Tác giả: IT Sectr Đã đăng: 2026-03-16 Thời gian đọc: 8 phút

RxJava là một thư viện lập trình phản ứng cho Java và Android, triển khai mẫu Observer thông qua Observable và Observer. Theo ReactiveX GitHub, 2026, RxJava cho phép xử lý các luồng dữ liệu và sự kiện bất đồng bộ bằng chuỗi toán tử. Đơn vị cơ bản là Observable, phát dữ liệu đến Observer thông qua một chuỗi các phép biến đổi. RxJava 3 là phiên bản ổn định hiện tại với hỗ trợ Java 8 lambda, Reactive Streams và tích hợp Android qua RxAndroid.

Những điểm chính

  • RxJava — triển khai Java của ReactiveX để xử lý luồng dữ liệu bất đồng bộ
  • Observable — nguồn dữ liệu phát các phần tử đến Observer
  • Observer — người đăng ký nhận thông báo onNext, onError và onComplete
  • Toán tử — chuỗi hàm để biến đổi, lọc và kết hợp các luồng
  • Schedulers — thành phần quản lý luồng thực thi của Observable và Observer

RxJava và ReactiveX là gì

RxJava là triển khai Java của đặc tả ReactiveX, một thư viện cho lập trình bất đồng bộ sử dụng các luồng có thể quan sát (Observable). RxJava 2 được phát hành vào năm 2016 với hỗ trợ Reactive Streams (Flowable) và phân chia thành rx.Observable và io.reactivex.Observable. RxJava 3 (2019) là phiên bản chính hiện tại với khả năng tương thích ngược với RxJava 2.

Ý tưởng cốt lõi của RxJava là mọi thứ đều là một luồng: luồng dữ liệu, luồng sự kiện, luồng trạng thái. Bất kỳ hoạt động bất đồng bộ nào cũng có thể được biểu diễn dưới dạng Observable phát dữ liệu, lỗi hoặc tín hiệu hoàn thành. Observer đăng ký Observable và nhận thông báo theo thời gian thực.

Theo Badoo (2024), trước khi chuyển sang coroutines, 76% ứng dụng Android trong top 200 của Google Play sử dụng RxJava cho các hoạt động bất đồng bộ. Tỷ lệ hiện đang giảm dần để chuyển sang coroutines, nhưng RxJava vẫn tồn tại trong mã sản xuất của hàng nghìn ứng dụng và được coi là một công nghệ trưởng thành, đã được kiểm chứng. ReactiveX là một đặc tả đa nền tảng cũng được triển khai cho JavaScript (RxJS), .NET (Rx.NET), Swift (RxSwift) và các ngôn ngữ khác.

Mẫu Observer trong RxJava

ReactiveX mở rộng mẫu Observer cổ điển với hai cơ chế: chuỗi toán tửquản lý luồng dựa trên Scheduler. Observable không bắt đầu phát dữ liệu cho đến khi Observer đăng ký (đánh giá lười). Điều này cho phép xây dựng một đường ống dữ liệu chỉ kích hoạt khi có đăng ký.

Các loại Observable: Observable, Flowable, Single, Maybe, Completable

Observable — loại cơ bản phát 0..N phần tử với onError hoặc onComplete. Phù hợp cho các luồng dữ liệu không giới hạn — ví dụ: sự kiện nhấp chuột hoặc cập nhật vị trí địa lý. Observable không hỗ trợ backpressure.

Flowable — phiên bản Reactive Streams của Observable với hỗ trợ backpressure. Được sử dụng khi nguồn dữ liệu có thể tạo ra phần tử nhanh hơn tốc độ xử lý của Observer. Flowable hỗ trợ các chiến lược BACKPRESSURE_BUFFER, DROP, LATEST và ERROR.

LoạiPhần tửBackpressureSử dụng
Observable0..NKhôngSự kiện UI, luồng nhỏ
Flowable0..NDữ liệu lớn, thời gian thực
Single1 (onSuccess/onError)Phản hồi đơn (mạng)
Maybe0..1Giá trị tùy chọn (bộ nhớ đệm)
Completable0 (onComplete/onError)Hoạt động không dữ liệu (ghi)

Single, Maybe và Completable

Single phát chính xác một phần tử hoặc lỗi — lý tưởng cho các yêu cầu mạng. Maybe phát 0 hoặc 1 phần tử, phù hợp cho bộ nhớ đệm nơi dữ liệu có thể không tồn tại. Completable chỉ phát onComplete hoặc onError, không có dữ liệu, thuận tiện cho các hoạt động ghi hoặc xóa. Các loại này đơn giản hóa API bằng cách thu hẹp hợp đồng thành một trường hợp cụ thể. Retrofit (một trình khách HTTP phổ biến cho Android) hỗ trợ trực tiếp cả năm loại RxJava, cho phép chọn loại trả về phù hợp nhất cho mỗi điểm cuối mà không cần mã mẫu thừa.

Toán tử RxJava: biến đổi và lọc luồng

Toán tử là các hàm biến đổi một Observable thành một Observable khác. Một chuỗi toán tử mô tả đường ống dữ liệu: mỗi toán tử nhận luồng từ toán tử trước, biến đổi nó và chuyển cho toán tử tiếp theo. RxJava chứa hơn 200 toán tử được nhóm thành các danh mục.

  • map — biến đổi mỗi phần tử (Integer → String)
  • flatMap — biến đổi một phần tử thành Observable và hợp nhất tất cả thành một luồng
  • filter — chỉ cho qua các phần tử thỏa mãn điều kiện
  • zip — kết hợp các phần tử từ N Observable theo chỉ mục
  • merge — hợp nhất nhiều Observable thành một, giữ nguyên thứ tự thời gian
  • debounce — chỉ phát phần tử nếu một khoảng thời gian xác định trôi qua mà không có phát nào khác

flatMap là một trong những toán tử mạnh nhất của RxJava. Nó cho phép thực thi một yêu cầu bất đồng bộ cho mỗi phần tử và thu thập kết quả vào một luồng chung. Ví dụ, flatMap được sử dụng để tải chi tiết từ danh sách ID: mỗi ID → yêu cầu mạng → hợp nhất kết quả. Không giống như map chỉ đơn giản biến đổi một phần tử, flatMap có thể phát nhiều phần tử hoặc chuyển sang một Observable khác, khiến nó trở thành nền tảng để xây dựng các đường ống bất đồng bộ.

Xử lý lỗi bằng toán tử

onErrorResumeNext — chuyển sang Observable dự phòng khi có lỗi. retry — đăng ký lại N lần khi có lỗi. onErrorReturn — trả về giá trị mặc định thay vì lỗi. doOnError — thực hiện tác dụng phụ khi có lỗi mà không thay đổi luồng (ghi nhật ký hoặc phân tích). Kết hợp các toán tử này cho phép xây dựng các đường ống mạnh mẽ với chiến lược xử lý lỗi rõ ràng mà không cần try/catch thủ công.

Schedulers: quản lý luồng trong RxJava

Schedulers xác định Observable và Observer thực thi trên luồng nào. subscribeOn đặt luồng cho nguồn, observeOn đặt luồng cho Observer và các toán tử tiếp theo. Sự tách biệt này là một lợi thế chính của RxJava: nguồn trên luồng IO, xử lý trên computation, UI trên luồng chính.

Các Schedulers chính: Schedulers.io() — cho các hoạt động I/O (mạng, đĩa), nhóm không giới hạn. Schedulers.computation() — cho tính toán, nhóm cố định theo số lõi. Schedulers.newThread() — một luồng mới cho mỗi tác vụ. AndroidSchedulers.mainThread() — luồng chính Android (RxAndroid). Cũng có Schedulers.trampoline() để thực thi tác vụ trong luồng hiện tại với hàng đợi FIFO, hữu ích cho kiểm thử.

Theo Google (2025), sử dụng đúng Schedulers là phần khó nhất của RxJava đối với người mới bắt đầu. Một lỗi điển hình là gọi subscribeOn sau observeOn, điều này không ảnh hưởng đến nguồn. subscribeOn phải là đầu tiên trong chuỗi cho nguồn, observeOn trước khi đăng ký UI. Quy tắc: subscribeOn chỉ ảnh hưởng đến thượng nguồn (nguồn), observeOn chuyển hạ nguồn (người đăng ký và tất cả toán tử sau nó).

Ví dụ mã RxJava trong Android

Hãy xem xét ba kịch bản: yêu cầu mạng với Single, yêu cầu song song với zip và debounce cho trường tìm kiếm với debounce.

Yêu cầu mạng với Single

Single hoàn hảo cho các yêu cầu Retrofit: một yêu cầu — một phản hồi. Đăng ký trên luồng chính để cập nhật UI.

java
api.getUser(id)
    .subscribeOn(Schedulers.io())
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(new SingleObserver<User>() {
        @Override
        public void onSuccess(User user) { showUser(user); }
        @Override
        public void onError(Throwable e) { showError(e); }
    })

Yêu cầu song song với zip

zip kết hợp kết quả của hai Single độc lập thành một. Chúng thực thi song song, kết quả được tạo ra sau khi cả hai hoàn thành.

java
Single.zip(
    api.getProfile(),
    api.getSettings(),
    (profile, settings) -> new Dashboard(profile, settings)
)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe(dashboard -> showDashboard(dashboard), e -> logError(e))

Debounce cho trường tìm kiếm

debounce bỏ qua các thay đổi văn bản nhanh và chỉ gửi yêu cầu sau khi tạm dừng 400 ms. distinctUntilChanged hủy yêu cầu nếu văn bản không thay đổi.

java
RxTextView.textChanges(searchView)
    .debounce(400, TimeUnit.MILLISECONDS)
    .filter(text -> text.length() >= 3)
    .distinctUntilChanged()
    .switchMap(query -> api.search(query))
    .observeOn(AndroidSchedulers.mainThread())
    .subscribe(results -> showResults(results))

RxJava so với Kotlin Coroutines: so sánh cách tiếp cận

RxJavaKotlin Coroutines giải quyết cùng một vấn đề — lập trình bất đồng bộ — nhưng với các cách tiếp cận khác nhau về cơ bản. RxJava được xây dựng trên mẫu Observer và dựa trên push: nguồn gửi dữ liệu, Observer phản ứng. Coroutines dựa trên pull: mã yêu cầu dữ liệu tuần tự thông qua await.

  • RxJava — phản ứng, luồng dữ liệu, >200 toán tử, push-based, đường cong học tập dốc
  • Coroutines — tuần tự, suspend/await, ~40 hàm, pull-based, cú pháp đơn giản
  • RxJava — trưởng thành (2016), hệ sinh thái lớn, nhưng đường cong học tập dốc
  • Coroutines — hiện đại (2018), lựa chọn ưa thích của Google cho mã mới
  • RxJava — backpressure tích hợp qua Flowable, chiến lược đệm được kiểm chứng
  • Coroutines — Flow với backpressure gần đây, nhưng đang được JetBrains phát triển tích cực

Theo Google I/O 2024, Kotlin Coroutines là cách tiếp cận được khuyến nghị cho mã bất đồng bộ mới trong Android. RxJava vẫn được hỗ trợ cho các dự án hiện có. Google cung cấp các thư viện cầu nối (kotlinx-coroutines-rx3) để di chuyển dần dần. AndroidX (LiveData, Room, Paging 3) hỗ trợ cả hai cách tiếp cận, cho phép sử dụng RxJava trong các mô-đun cũ và coroutines trong các mô-đun mới mà không xung đột phụ thuộc.

Chiến lược di chuyển từ RxJava sang Coroutines

Chuyển đổi dần dần: mỗi thành phần mới được viết bằng coroutines, mã RxJava cũ không bị sửa đổi. RxJava → coroutines qua awaitSingle() hoặc awaitFirst(). Coroutines → RxJava qua future() hoặc asFlowable(). Di chuyển hoàn toàn mất 6–18 tháng cho các dự án lớn.

Câu hỏi thường gặp

Observable khác Flowable như thế nào?

Observable không hỗ trợ backpressure — nếu nguồn tạo dữ liệu nhanh hơn tốc độ xử lý của trình xử lý, MissingBackpressureException sẽ xảy ra. Flowable hỗ trợ backpressure Reactive Streams với các chiến lược đệm có thể cấu hình.

subscribeOn và observeOn là gì?

subscribeOn đặt Scheduler để thực thi Observable nguồn. observeOn đặt Scheduler cho Observer và tất cả toán tử tiếp theo trong chuỗi. subscribeOn ảnh hưởng đến thượng nguồn, observeOn ảnh hưởng đến hạ nguồn.

Tôi có nên chuyển từ RxJava sang coroutines không?

Đối với dự án mới — có, Google khuyến nghị coroutines. Đối với dự án hiện có — di chuyển dần dần qua kotlinx-coroutines-rx3. RxJava vẫn ổn định và được hỗ trợ cho mã cũ.

Làm thế nào để xử lý lỗi trong RxJava?

Thông qua toán tử: onErrorReturn (giá trị mặc định), onErrorResumeNext (Observable dự phòng), retry (thử lại N lần). Hoặc qua Observer.onError() để hiển thị cho người dùng.

CompositeDisposable là gì?

CompositeDisposable là một vùng chứa để quản lý nhiều đăng ký. Khi dispose() được gọi, tất cả các đăng ký đã thêm đều bị hủy. Nó được sử dụng trong Activity/Fragment để hủy tất cả yêu cầu khi màn hình bị hủy.

Tổng kết

  • RxJava — thư viện lập trình phản ứng cho Java và Android dựa trên mẫu Observer
  • Observable/Flowable — nguồn dữ liệu có và không hỗ trợ backpressure
  • Single, Maybe, Completable — loại chuyên biệt cho 1, 0..1 và 0 phần tử
  • Toán tử (map, flatMap, zip, filter) — chuỗi biến đổi với hơn 200 hàm
  • Schedulers — subscribeOn cho nguồn và observeOn cho người tiêu dùng
  • RxJava vs Coroutines — Google khuyến nghị coroutines cho mã mới, RxJava cho mã cũ
  • CompositeDisposable — quản lý đăng ký an toàn với hủy bỏ khi màn hình bị hủy

Chúng tôi sẽ phát triển ứng dụng di động chìa khóa trao tay

IT Sectr tạo các ứng dụng iOS và Android cho các công ty khởi nghiệp và doanh nghiệp từ năm 2017. Chúng tôi sẽ tư vấn và đề xuất giải pháp tốt nhất cho bạn.

Thảo luận dự án

Đọc thêm