RxJava คือไลบรารีการเขียนโปรแกรมแบบรีแอกทีฟสำหรับ JVM ที่ implements สตรีมข้อมูลแบบอะซิงโครนัสผ่านรูปแบบ Observable พร้อมด้วยโอเปอเรเตอร์เชิงฟังก์ชันสำหรับการแปลงค่า มัน移植แนวคิดของ ReactiveX ไปยัง Java และ Kotlin โดยมี API แบบรวมสำหรับทำงานกับคำขอเครือข่าย ฐานข้อมูล อีเวนต์ UI และงานเบื้องหลัง ตามข้อมูลจาก ReactiveX, 2025 ไลบรารีนี้ถูกใช้ในกว่า 120,000 โปรเจกต์บน GitHub และเป็นมาตรฐานการเขียนโปรแกรมแบบรีแอกทีฟสำหรับ Android จนกระทั่ง Kotlin Flow ปรากฏ RxJava แทนที่ AsyncTask, Loader และ callback ด้วยเชนการประมวลผลข้อมูลเดียว
ประเด็นสำคัญ
RxJava คือการ implement ไลบรารี ReactiveX (Reactive Extensions) สำหรับ Java Virtual Machine เวอร์ชันแรกของ RxJava ถูกเผยแพร่โดย Netflix ในปี 2013 เพื่อจัดการการเรียกแบบอะซิงโครนัสในแอปพลิเคชันฝั่งเซิร์ฟเวอร์ ในช่วงเวลาที่สร้าง ทางเลือกหลักใน Java คือ Future และ Callback — ทั้งสองวิธีนำไปสู่ callback-hell และการจัดการเธรดที่ซับซ้อน RxJava ได้แนะนำการประกอบการดำเนินการแบบอะซิงโครนัสผ่าน Observable ด้วยเชนของโอเปอเรเตอร์เชิงฟังก์ชัน
สถาปัตยกรรมของ RxJava ขึ้นอยู่กับข้อกำหนด Reactive Streams — มาตรฐานสำหรับการประมวลผลสตรีมแบบอะซิงโครนัสด้วย backpressure ที่ไม่บล็อก ข้อกำหนดกำหนดอินเทอร์เฟสสี่แบบ: Publisher, Subscriber, Subscription และ Processor RxJava 2+ implement Reactive Streams อย่างสมบูรณ์ผ่านประเภท Flowable โดยปฏิบัติตามสัญญา backpressure ซึ่งแตกต่างจาก RxJava 1 Observable ใน RxJava 2 ไม่สนับสนุน backpressure — มันมีไว้สำหรับสตรีมที่มีอีเวนต์จำนวนน้อยหรืออีเวนต์ UI
จากการสำรวจของ JetBrains, 2025 RxJava อยู่ใน 3 อันดับแรกของไลบรารีสำหรับการพัฒนา Android กรณีการใช้งานหลักรวมถึง: การจัดการคำขอเครือข่ายผ่าน Retrofit (รวมกับ RxJava ผ่าน CallAdapter), การทำงานกับ Room (คำสั่งรีแอกทีฟส่งคืน Flowable หรือ Maybe), แอนิเมชันและอีเวนต์ UI ผ่าน RxBinding และการค้นหาแบบ debounce เมื่อป้อนข้อความ ทุกสถานการณ์เหล่านี้ใช้รูปแบบเชนร่วมกัน: แหล่งที่มา (Observable) → การแปลง (โอเปอเรเตอร์) → การสมัครสมาชิก (subscribe)
RxJava 1 (2013) วางรากฐานด้วย Observable และโอเปอเรเตอร์ แต่ประสบปัญหา backpressure — ในสตรีมที่เร็ว ข้อมูลสะสมในหน่วยความจำ ทำให้เกิด OutOfMemoryError RxJava 2 (2016) แก้ไขสถาปัตยกรรมโดยแยก Observable (ไม่มี backpressure) และ Flowable (มี backpressure) RxJava 3 (2020) เพิ่มการสนับสนุน Java 8 Stream API โอเปอเรเตอร์เพิ่มเติม และประสิทธิภาพการสมัครสมาชิกที่ดีขึ้น ปัจจุบัน RxJava 3 เป็นเวอร์ชันที่แนะนำสำหรับโปรเจกต์ใหม่
RxJava มีแหล่งที่มารีแอกทีฟหลักห้าประเภท แต่ละประเภทออกแบบมาสำหรับสถานการณ์เฉพาะ Observable และ Flowable ปล่อยค่าหลายค่า Single ปล่อยหนึ่งค่าหรือข้อผิดพลาด Completable ปล่อยเฉพาะการเสร็จสิ้นโดยไม่มีข้อมูล และ Maybe ปล่อยหนึ่งค่า ศูนย์ หรือข้อผิดพลาด การเลือกประเภทที่ถูกต้องช่วยลดปริมาณโค้ดและทำให้เชนอธิบายตนเองได้
| ประเภท | จำนวนอีเวนต์ | Backpressure | สถานการณ์ |
|---|---|---|---|
| Observable | 0..N จากนั้นเสร็จสิ้น | ไม่ | อีเวนต์ UI, สตรีมสั้น |
| Flowable | 0..N จากนั้นเสร็จสิ้น | ใช่ | การตอบสนองเครือข่าย, สตรีม DB |
| Single | 1 พอดีหรือข้อผิดพลาด | ไม่ | คำขอ HTTP, อ่านหนึ่งระเบียน |
| Completable | 0 (เฉพาะการเสร็จสิ้น) | ไม่ | การเขียน DB, การส่งอีเวนต์ |
| Maybe | 0, 1 หรือข้อผิดพลาด | ไม่ | แคช: มีค่าหรือไม่มี |
Flowable เป็นประเภทที่ยืดหยุ่นที่สุดสำหรับการทำงานกับสตรีมข้อมูลขนาดใหญ่ มัน implement Publisher ของ Reactive Streams พร้อมการสนับสนุน backpressure: consumer สามารถขอจำนวนองค์ประกอบที่เฉพาะเจาะจงผ่าน Subscription.request(n) ซึ่งป้องกันบัฟเฟอร์ล้นเมื่อความเร็วของ producer และ consumer ไม่ตรงกัน ถ้า backpressure ไม่สำคัญ ให้ใช้ Observable — มันมีโอเวอร์เฮดน้อยกว่าเนื่องจากไม่มีกลไกการขอ
Single เป็นตัวเลือกที่ดีที่สุดสำหรับคำขอ HTTP Retrofit 2 กับ RxJava CallAdapter ส่งคืน Single<ResponseBody> สำหรับแต่ละคำขอ Single รับประกันการเรียก onSuccess หรือ onError หนึ่งครั้งพอดี ซึ่งสอดคล้องกับความหมายของคำขอ HTTP — หนึ่งการตอบสนองหรือหนึ่งข้อผิดพลาด Completable ใช้สำหรับการดำเนินการเขียนที่ไม่ส่งคืนข้อมูล: insert, update, delete Maybe สะดวกสำหรับการตรวจสอบแคช — มันอาจส่งคืนค่าหรือไม่ก็ได้
// ตัวอย่างการใช้ Single สำหรับคำขอ HTTP
interface ApiService {
@GET("users/{id}")
fun getUser(@Path("id") userId: Int): Single<User>
}
// การสมัครสมาชิกด้วยการประมวลผลบนเธรดหลัก
apiService.getUser(42)
.subscribeOn(Schedulers.io())
.observeOn(AndroidSchedulers.mainThread())
.subscribe({ user ->
textView.text = user.name
}, { error ->
Log.e("API", "Error: ${error.message}")
})
.addTo(compositeDisposable)
โอเปอเรเตอร์ ใน RxJava คือฟังก์ชันลำดับสูงที่รับแหล่งที่มารีแอกทีฟหนึ่งและส่งคืนอีกแหล่งหนึ่ง โดยแปลงสตรีมข้อมูล RxJava 3 มีโอเปอเรเตอร์มากกว่า 400 ตัว แบ่งเป็นหมวดหมู่: การแปลง การกรอง การรวม การจัดการข้อผิดพลาด และการจัดการเวลา โอเปอเรเตอร์แต่ละตัวเฉื่อย — เชนถูกสร้างขึ้นที่การประกาศและดำเนินการที่การสมัครสมาชิก
map คือโอเปอเรเตอร์พื้นฐานที่แปลงแต่ละค่าผ่านฟังก์ชัน flatMap รับฟังก์ชันที่ส่งคืน Observable สำหรับแต่ละองค์ประกอบ และทำให้ผลลัพธ์เรียบเป็นสตรีมเดียว switchMap คล้ายกับ flatMap แต่เมื่อองค์ประกอบใหม่มาถึง มันยกเลิกการสมัครสมาชิกจาก Observable ก่อนหน้า concatMap รักษาลำดับองค์ประกอบ — แตกต่างจาก flatMap มันสมัครสมาชิกตามลำดับไปยัง Observable ที่ซ้อนกันแต่ละตัว
// การแยกวิเคราะห์ JSON ด้วยการแปลงและการกรอง
apiService.getUsers()
.flatMap { users ->
Observable.fromIterable(users)
}
.filter { user ->
user.age >= 18
}
.map { user ->
UserDto(user.name, user.age)
}
.toList()
.subscribeOn(Schedulers.computation())
.observeOn(AndroidSchedulers.mainThread())
.subscribe({ adapter.submitList(it) },
{ Log.e("ข้อผิดพลาด", it.message) })
การรวมสตรีม เป็นพื้นที่ที่ RxJava มีความโดดเด่นเป็นพิเศษ zip รวมองค์ประกอบจาก Observable หลายตัวเป็นคู่ตามดัชนี: ตัวแรกกับตัวแรก ตัวที่สองกับตัวที่สอง combineLatest ปล่อยค่าใหม่เมื่อสตรีมใด ๆ เปลี่ยนแปลง โดยรวมค่าล่าสุดจากทุกสตรีม merge รวม Observable หลายตัวเป็นหนึ่งเดียว โดยรักษาลำดับการมาถึงของอีเวนต์ concat สมัครสมาชิกตามลำดับไปยัง Observable แต่ละตัวและส่งต่ออีเวนต์ทั้งหมดก่อนที่จะไปยังตัวถัดไป
การจัดการเวลา รวมถึง debounce (รอการหยุดชั่วคราวในสตรีมก่อนปล่อย), throttleFirst (ปล่อยอีเวนต์แรก ไม่สนใจส่วนที่เหลือภายในหน้าต่าง), timeout (ข้อผิดพลาดหากไม่มีอีเวนต์มาถึงภายในช่วงเวลา) การค้นหาแบบ debounce เมื่อป้อนข้อความเป็นสถานการณ์ที่พบบ่อยที่สุด: searchObservable.debounce(300, MILLISECONDS).distinctUntilChanged() ป้องกันคำขอที่ไม่จำเป็นระหว่างการพิมพ์เร็ว
| หมวดหมู่ | โอเปอเรเตอร์ | พฤติกรรม |
|---|---|---|
| การแปลง | map / flatMap / switchMap | แปลงค่าหรือสตรีม |
| การกรอง | filter / distinct / take | เลือกค่าตามเงื่อนไข |
| การรวม | zip / combineLatest / merge | รวม 2+ สตรีม |
| ข้อผิดพลาด | onErrorResumeNext / retry | กู้คืนจากความล้มเหลว |
| ยูทิลิตี้ | delay / timeout / debounce | การจัดการเวลาในสตรีม |
Scheduler ใน RxJava คือสิ่งที่เป็นนามธรรมเหนือพูลเธรด ไลบรารีมี Scheduler ในตัวห้าตัว: Schedulers.io() สำหรับการดำเนินการ I/O (เครือข่าย, ไฟล์), Schedulers.computation() สำหรับงานที่ใช้ CPU มาก, Schedulers.newThread() สำหรับเธรดใหม่ทุกครั้ง, Schedulers.single() สำหรับการดำเนินการเธรดเดียว และ Schedulers.trampoline() สำหรับการดำเนินการทันทีในเธรดปัจจุบัน
subscribeOn กำหนดว่า Scheduler ใดดำเนินการ Observable ต้นทาง ถ้ามี subscribeOn หลายตัวในเชน ตัวที่ใกล้ต้นทางที่สุดจะมีความสำคัญ observeOn สลับ downstream ไปยัง Scheduler ที่ระบุ — การใช้ observeOn แต่ละครั้งเปลี่ยนเธรดสำหรับโอเปอเรเตอร์ที่ตามมา รูปแบบ Android ทั่วไป: subscribeOn(Schedulers.io()) สำหรับการดำเนินการเครือข่าย, observeOn(AndroidSchedulers.mainThread()) สำหรับการอัปเดต UI
// การประมวลผลแบบหลายเธรดด้วยการสลับบริบท
Observable.fromCallable(() -> database.getItems())
.subscribeOn(Schedulers.io()) // DB บน io
.map(items -> processItems(items)) // การแปลงบน io
.observeOn(Schedulers.computation()) // สลับไป computation
.map(processed -> compressImages(processed))
.observeOn(AndroidSchedulers.mainThread())
.subscribe(result -> ui.showResult(result))
AndroidSchedulers.mainThread() คือ Scheduler จากไลบรารี RxAndroid ที่ดำเนินการโค้ดบนเธรดหลักของ Android มันจำเป็นสำหรับการอัปเดต UI ใด ๆ ในเชนรีแอกทีฟ ไลบรารีใช้ Handler ภายในและรับประกันการดำเนินการบนเธรด UI แม้ภายใต้โหลดสูง สำหรับการดำเนินการเบื้องหลัง Schedulers.io() รองรับพูลเธรดไม่จำกัดและเหมาะสำหรับการดำเนินการที่บล็อกใด ๆ Schedulers.computation() ใช้พูลคงที่เท่ากับจำนวนแกน CPU
RxJava ใน Android ใช้สำหรับสามสถานการณ์หลัก: คำสั่งรีแอกทีฟไปยัง Room, การรวมกับ Retrofit และการผูก UI แบบรีแอกทีฟผ่าน RxBinding แต่ละสถานการณ์มีชุดประเภทของตัวเอง: Room ส่งคืน Flowable สำหรับคำสั่งที่สังเกตได้, Retrofit ส่งคืน Single สำหรับคำขอ HTTP, RxBinding ส่งคืน Observable สำหรับอีเวนต์ UI
Room คือไลบรารีการคงอยู่ของข้อมูลจาก Google เริ่มต้นจาก Room 2.1 ฐานข้อมูลสนับสนุนประเภทส่งคืนรีแอกทีฟ: Flowable และ Observable เมื่อระเบียนใด ๆ ในตารางเปลี่ยนแปลง Room จะส่งค่าใหม่ไปยังสตรีมโดยอัตโนมัติ นักพัฒนาสมัครสมาชิก Flowable ใน ViewModel และรับข้อมูลที่อัปเดตโดยไม่ต้องสืบค้นด้วยตนเองทุกครั้งที่เปลี่ยนแปลง
// Room DAO ด้วยคำสั่งรีแอกทีฟ
@Dao
interface UserDao {
@Query("SELECT * FROM users WHERE id = :id")
fun getUserById(@Param("id") userId: Int): Flowable<User>
@Insert
fun insertUser(user: User): Completable
}
// ViewModel — องค์ประกอบ Room + Network
class UserViewModel(private val dao: UserDao) : ViewModel() {
val users: Flowable<List<User>> = dao.getAllUsers()
.subscribeOn(Schedulers.io())
}
รูปแบบ MVVM + RxJava สร้างขึ้นบนหลักการที่ ViewModel ไม่มีการอ้างอิงไปยัง View ViewModel เผยแพร่แหล่งที่มารีแอกทีฟ (Flowable, LiveData ผ่าน Transformations) และ Activity หรือ Fragment สมัครสมาชิก สิ่งนี้ให้ความสามารถในการทดสอบ: ViewModel ถูกทดสอบโดยไม่มี UI โดยการแทนที่ Scheduler ผ่าน RxJavaPlugins.setComputationScheduler CompositeDisposable ใน ViewModel จัดการวงจรชีวิตของการสมัครสมาชิก — เมื่อ onCleared() การสมัครสมาชิกทั้งหมดจะถูกยกเลิก
Kotlin Flow คือการ implement สตรีมเย็นแบบพื้นเมืองใน Kotlin ซึ่งสร้างใน coroutines และแนะนำใน Kotlin 1.3 Flow แก้ปัญหาเดียวกันกับ RxJava แต่มีความแตกต่างพื้นฐาน: การสนับสนุน coroutine ในตัว (ฟังก์ชัน suspend), การยกเลิกผ่าน coroutine cancellation และไม่มีปัญหา backpressure — Flow ใช้ suspend แทนการบัฟเฟอร์ Flow เป็นส่วนหนึ่งของไลบรารีมาตรฐาน Kotlin โดยไม่ต้องพึ่งพาเพิ่มเติม
RxJava ยังคงเป็นตัวเลือกที่ต้องการสำหรับโปรเจกต์ Java, โปรเจกต์ที่สนับสนุน Java 7-8 และฐานโค้ด RxJava ที่มีอยู่ ระบบนิเวศของ RxJava รวยกว่าอย่างมาก: โอเปอเรเตอร์มากกว่า 400 ตัวเทียบกับประมาณ 50 ใน Flow, การรวมกับ Retrofit ผ่าน CallAdapter ในตัว, การสนับสนุน backpressure ผ่าน Flowable และ RxBinding, RxPermissions, RxLocation สำหรับ Android Kotlin Flow กำลังตามมาอย่างรวดเร็ว แต่ความยืดหยุ่นของ RxJava ในสถานการณ์การรวมสตรีมที่ซับซ้อนยังคงสูงกว่า
| ลักษณะ | RxJava | Kotlin Flow |
|---|---|---|
| ภาษา | Java / Kotlin | Kotlin เท่านั้น |
| การยกเลิก | Disposable / CompositeDisposable | Coroutine cancellation |
| Backpressure | Flowable (กลยุทธ์ BUFFER, DROP, LATEST) | ผ่าน conflate / buffer |
| โอเปอเรเตอร์ | 400+ | ~50 (ขยายได้) |
| การรวม Room | Flowable, Observable | Flow, StateFlow |
| ViewModel | CompositeDisposable | viewModelScope + Flow |
คำถามที่พบบ่อย
Observable ไม่สนับสนุน backpressure — ถ้า producer เร็วกว่า consumer อีเวนต์จะสะสมในหน่วยความจำ Flowable implement Reactive Streams ด้วย backpressure ผ่าน Subscription.request() ป้องกันบัฟเฟอร์ล้นเมื่อความเร็วไม่ตรงกัน
Single ใช้สำหรับการดำเนินการที่ส่งคืนหนึ่งค่าหรือข้อผิดพลาดพอดี: คำขอ HTTP, การอ่านหนึ่งระเบียนจาก DB, การคำนวณผลลัพธ์ Single มีความหมายตรงกับ Future และลดโค้ดโดยลบ onComplete ที่ไม่ได้ใช้
เมธอด dispose() บน Disposable ยกเลิกการสมัครสมาชิก สำหรับการจัดการกลุ่มใช้ CompositeDisposable — มันรวบรวม Disposable ทั้งหมดและยกเลิกพร้อมกันเมื่อเรียก clear() ตำแหน่งทั่วไปคือ onCleared() ใน ViewModel หรือ onPause() ใน Activity
flatMap สมัครสมาชิก Observable ที่ซ้อนกันทั้งหมดและรวมอีเวนต์ของพวกมันในลำดับใดก็ได้ switchMap เมื่อองค์ประกอบใหม่มาถึง ยกเลิกการสมัครสมาชิกจาก Observable ก่อนหน้าและสมัครสมาชิกใหม่ switchMap ใช้ในการค้นหา — แต่ละคำขอใหม่ยกเลิกคำขอก่อนหน้า
สำหรับโปรเจกต์ Kotlin ใหม่ Flow เหมาะสมกว่าเนื่องจากรวมกับ coroutines และขนาดเล็กกว่า สำหรับโปรเจกต์ RxJava ที่มีอยู่ การย้ายจะสมเหตุสมผลก็ต่อเมื่อฐานโค้ดทั้งหมดกำลังย้ายไป coroutines — การใช้ทั้งสองไลบรารีระหว่างทางทำให้สถาปัตยกรรมซับซ้อน
สรุป
เราจะพัฒนาแอปพลิเคชันบนมือถือแบบครบวงจร
IT Sectr สร้างแอปพลิเคชัน iOS และ Android สำหรับสตาร์ทอัพและธุรกิจตั้งแต่ปี 2017 เราจะให้คำแนะนำและเสนอวิธีแก้ปัญหาที่ดีที่สุดแก่คุณ