RxJava: แก่นแท้ ส่วนประกอบ และการเขียนโปรแกรมแบบรีแอกทีฟ

ผู้แต่ง: IT Sectr เผยแพร่เมื่อ: 2026-05-03 เวลาอ่าน: 10 นาที

RxJava คือไลบรารีการเขียนโปรแกรมแบบรีแอกทีฟสำหรับ JVM ที่ implements สตรีมข้อมูลแบบอะซิงโครนัสผ่านรูปแบบ Observable พร้อมด้วยโอเปอเรเตอร์เชิงฟังก์ชันสำหรับการแปลงค่า มัน移植แนวคิดของ ReactiveX ไปยัง Java และ Kotlin โดยมี API แบบรวมสำหรับทำงานกับคำขอเครือข่าย ฐานข้อมูล อีเวนต์ UI และงานเบื้องหลัง ตามข้อมูลจาก ReactiveX, 2025 ไลบรารีนี้ถูกใช้ในกว่า 120,000 โปรเจกต์บน GitHub และเป็นมาตรฐานการเขียนโปรแกรมแบบรีแอกทีฟสำหรับ Android จนกระทั่ง Kotlin Flow ปรากฏ RxJava แทนที่ AsyncTask, Loader และ callback ด้วยเชนการประมวลผลข้อมูลเดียว

ประเด็นสำคัญ

  • RxJava คือการimplement ReactiveX สำหรับ Java/Kotlin ด้วยประเภท Observable, Flowable, Single, Completable และ Maybe
  • Observable แสดงถึงสตรีมข้อมูลที่มีการจัดการ backpressure ผ่าน Flowable เมื่อสมัครสมาชิกบน consumer ที่ช้า
  • โอเปอเรเตอร์ map, flatMap, switchMap, zip และ combineLatest แปลงและรวมสตรีมแบบอะซิงโครนัสโดยไม่บล็อก
  • Scheduler — Schedulers.io(), computation(), mainThread() จัดการว่าเธรดใดทำงานและสมัครสมาชิก
  • RxAndroid เพิ่ม AndroidSchedulers.mainThread() สำหรับอัปเดต UI จากเชนรีแอกทีฟ

RxJava คืออะไร?

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

RxJava 1 (2013) วางรากฐานด้วย Observable และโอเปอเรเตอร์ แต่ประสบปัญหา backpressure — ในสตรีมที่เร็ว ข้อมูลสะสมในหน่วยความจำ ทำให้เกิด OutOfMemoryError RxJava 2 (2016) แก้ไขสถาปัตยกรรมโดยแยก Observable (ไม่มี backpressure) และ Flowable (มี backpressure) RxJava 3 (2020) เพิ่มการสนับสนุน Java 8 Stream API โอเปอเรเตอร์เพิ่มเติม และประสิทธิภาพการสมัครสมาชิกที่ดีขึ้น ปัจจุบัน RxJava 3 เป็นเวอร์ชันที่แนะนำสำหรับโปรเจกต์ใหม่

ประเภทของสตรีมรีแอกทีฟใน RxJava

RxJava มีแหล่งที่มารีแอกทีฟหลักห้าประเภท แต่ละประเภทออกแบบมาสำหรับสถานการณ์เฉพาะ Observable และ Flowable ปล่อยค่าหลายค่า Single ปล่อยหนึ่งค่าหรือข้อผิดพลาด Completable ปล่อยเฉพาะการเสร็จสิ้นโดยไม่มีข้อมูล และ Maybe ปล่อยหนึ่งค่า ศูนย์ หรือข้อผิดพลาด การเลือกประเภทที่ถูกต้องช่วยลดปริมาณโค้ดและทำให้เชนอธิบายตนเองได้

ประเภทจำนวนอีเวนต์Backpressureสถานการณ์
Observable0..N จากนั้นเสร็จสิ้นไม่อีเวนต์ UI, สตรีมสั้น
Flowable0..N จากนั้นเสร็จสิ้นใช่การตอบสนองเครือข่าย, สตรีม DB
Single1 พอดีหรือข้อผิดพลาดไม่คำขอ HTTP, อ่านหนึ่งระเบียน
Completable0 (เฉพาะการเสร็จสิ้น)ไม่การเขียน DB, การส่งอีเวนต์
Maybe0, 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 สะดวกสำหรับการตรวจสอบแคช — มันอาจส่งคืนค่าหรือไม่ก็ได้

kotlin
// ตัวอย่างการใช้ 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 ที่ซ้อนกันแต่ละตัว

kotlin
// การแยกวิเคราะห์ 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 และการทำงานแบบหลายเธรด

Scheduler ใน RxJava คือสิ่งที่เป็นนามธรรมเหนือพูลเธรด ไลบรารีมี Scheduler ในตัวห้าตัว: Schedulers.io() สำหรับการดำเนินการ I/O (เครือข่าย, ไฟล์), Schedulers.computation() สำหรับงานที่ใช้ CPU มาก, Schedulers.newThread() สำหรับเธรดใหม่ทุกครั้ง, Schedulers.single() สำหรับการดำเนินการเธรดเดียว และ Schedulers.trampoline() สำหรับการดำเนินการทันทีในเธรดปัจจุบัน

subscribeOn และ observeOn

subscribeOn กำหนดว่า Scheduler ใดดำเนินการ Observable ต้นทาง ถ้ามี subscribeOn หลายตัวในเชน ตัวที่ใกล้ต้นทางที่สุดจะมีความสำคัญ observeOn สลับ downstream ไปยัง Scheduler ที่ระบุ — การใช้ observeOn แต่ละครั้งเปลี่ยนเธรดสำหรับโอเปอเรเตอร์ที่ตามมา รูปแบบ Android ทั่วไป: subscribeOn(Schedulers.io()) สำหรับการดำเนินการเครือข่าย, observeOn(AndroidSchedulers.mainThread()) สำหรับการอัปเดต UI

java
// การประมวลผลแบบหลายเธรดด้วยการสลับบริบท
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: การประยุกต์ใช้เชิงปฏิบัติ

RxJava ใน Android ใช้สำหรับสามสถานการณ์หลัก: คำสั่งรีแอกทีฟไปยัง Room, การรวมกับ Retrofit และการผูก UI แบบรีแอกทีฟผ่าน RxBinding แต่ละสถานการณ์มีชุดประเภทของตัวเอง: Room ส่งคืน Flowable สำหรับคำสั่งที่สังเกตได้, Retrofit ส่งคืน Single สำหรับคำขอ HTTP, RxBinding ส่งคืน Observable สำหรับอีเวนต์ UI

Room + RxJava

Room คือไลบรารีการคงอยู่ของข้อมูลจาก Google เริ่มต้นจาก Room 2.1 ฐานข้อมูลสนับสนุนประเภทส่งคืนรีแอกทีฟ: Flowable และ Observable เมื่อระเบียนใด ๆ ในตารางเปลี่ยนแปลง Room จะส่งค่าใหม่ไปยังสตรีมโดยอัตโนมัติ นักพัฒนาสมัครสมาชิก Flowable ใน ViewModel และรับข้อมูลที่อัปเดตโดยไม่ต้องสืบค้นด้วยตนเองทุกครั้งที่เปลี่ยนแปลง

kotlin
// 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() การสมัครสมาชิกทั้งหมดจะถูกยกเลิก

RxJava เปรียบเทียบกับ Kotlin Flow

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 ในสถานการณ์การรวมสตรีมที่ซับซ้อนยังคงสูงกว่า

ลักษณะRxJavaKotlin Flow
ภาษาJava / KotlinKotlin เท่านั้น
การยกเลิกDisposable / CompositeDisposableCoroutine cancellation
BackpressureFlowable (กลยุทธ์ BUFFER, DROP, LATEST)ผ่าน conflate / buffer
โอเปอเรเตอร์400+~50 (ขยายได้)
การรวม RoomFlowable, ObservableFlow, StateFlow
ViewModelCompositeDisposableviewModelScope + Flow

คำถามที่พบบ่อย

ความแตกต่างระหว่าง Observable และ Flowable ใน RxJava คืออะไร?

Observable ไม่สนับสนุน backpressure — ถ้า producer เร็วกว่า consumer อีเวนต์จะสะสมในหน่วยความจำ Flowable implement Reactive Streams ด้วย backpressure ผ่าน Subscription.request() ป้องกันบัฟเฟอร์ล้นเมื่อความเร็วไม่ตรงกัน

เมื่อใดควรใช้ Single แทน Observable?

Single ใช้สำหรับการดำเนินการที่ส่งคืนหนึ่งค่าหรือข้อผิดพลาดพอดี: คำขอ HTTP, การอ่านหนึ่งระเบียนจาก DB, การคำนวณผลลัพธ์ Single มีความหมายตรงกับ Future และลดโค้ดโดยลบ onComplete ที่ไม่ได้ใช้

จะยกเลิกการสมัครสมาชิกใน RxJava ได้อย่างไร?

เมธอด dispose() บน Disposable ยกเลิกการสมัครสมาชิก สำหรับการจัดการกลุ่มใช้ CompositeDisposable — มันรวบรวม Disposable ทั้งหมดและยกเลิกพร้อมกันเมื่อเรียก clear() ตำแหน่งทั่วไปคือ onCleared() ใน ViewModel หรือ onPause() ใน Activity

ความแตกต่างระหว่าง flatMap และ switchMap คืออะไร?

flatMap สมัครสมาชิก Observable ที่ซ้อนกันทั้งหมดและรวมอีเวนต์ของพวกมันในลำดับใดก็ได้ switchMap เมื่อองค์ประกอบใหม่มาถึง ยกเลิกการสมัครสมาชิกจาก Observable ก่อนหน้าและสมัครสมาชิกใหม่ switchMap ใช้ในการค้นหา — แต่ละคำขอใหม่ยกเลิกคำขอก่อนหน้า

ควรย้ายจาก RxJava ไป Kotlin Flow หรือไม่?

สำหรับโปรเจกต์ Kotlin ใหม่ Flow เหมาะสมกว่าเนื่องจากรวมกับ coroutines และขนาดเล็กกว่า สำหรับโปรเจกต์ RxJava ที่มีอยู่ การย้ายจะสมเหตุสมผลก็ต่อเมื่อฐานโค้ดทั้งหมดกำลังย้ายไป coroutines — การใช้ทั้งสองไลบรารีระหว่างทางทำให้สถาปัตยกรรมซับซ้อน

สรุป

  • RxJava คือไลบรารี ReactiveX สำหรับ JVM ด้วยประเภท Observable, Flowable, Single, Completable และ Maybe สำหรับสถานการณ์ต่าง ๆ
  • Flowable สนับสนุน backpressure ผ่าน Reactive Streams เพื่อป้องกันล้นเมื่อความเร็วไม่ตรงกัน
  • โอเปอเรเตอร์ map, flatMap, switchMap, zip, combineLatest, debounce ให้การประมวลผลสตรีมแบบประกาศ
  • Scheduler io(), computation(), mainThread() จัดการเธรดการดำเนินการโดยไม่บล็อก UI
  • RxAndroid รวม RxJava กับ Android โดยให้ AndroidSchedulers.mainThread() และทำให้การอัปเดต UI ง่ายขึ้น
  • Kotlin Flow เป็นทางเลือกพื้นเมืองที่มีการรวม coroutine แต่ RxJava รักษาข้อได้เปรียบในระบบนิเวศโอเปอเรเตอร์
  • MVVM + RxJava เป็นรูปแบบการพัฒนา Android มาตรฐานด้วย ViewModel ที่แยกจาก UI และการสมัครสมาชิกรีแอกทีฟ

เราจะพัฒนาแอปพลิเคชันบนมือถือแบบครบวงจร

IT Sectr สร้างแอปพลิเคชัน iOS และ Android สำหรับสตาร์ทอัพและธุรกิจตั้งแต่ปี 2017 เราจะให้คำแนะนำและเสนอวิธีแก้ปัญหาที่ดีที่สุดแก่คุณ

ปรึกษาโครงการ

อ่านเพิ่มเติม