RxJava: พื้นฐาน, ReactiveX และการทำงานกับสตรีมข้อมูล

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

RxJava เป็นไลบรารีการเขียนโปรแกรมเชิงรีแอคทีฟสำหรับ Java และ Android ที่ใช้รูปแบบ Observer ผ่าน Observable และ Observer ตามข้อมูลจาก ReactiveX GitHub, 2026 RxJava ช่วยให้สามารถจัดการสตรีมข้อมูลและเหตุการณ์แบบอะซิงโครนัสด้วยลูกโซ่ของโอเปอเรเตอร์ หน่วยพื้นฐานคือ Observable ซึ่งส่งข้อมูลไปยัง Observer ผ่านลูกโซ่ของการแปลง RxJava 3 เป็นเวอร์ชันเสถียรปัจจุบันที่รองรับ Java 8 lambda, Reactive Streams และการผสานรวมกับ Android ผ่าน RxAndroid

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

  • RxJava — การนำ ReactiveX มาใช้ใน Java สำหรับการประมวลผลสตรีมข้อมูลแบบอะซิงโครนัส
  • Observable — แหล่งข้อมูลที่ส่งองค์ประกอบไปยัง Observer
  • Observer — ผู้สมัครสมาชิกที่ได้รับการแจ้งเตือน onNext, onError และ onComplete
  • โอเปอเรเตอร์ — ลูกโซ่ของฟังก์ชันสำหรับการแปลง กรอง และรวมสตรีม
  • Schedulers — คอมโพเนนต์สำหรับจัดการเธรดการทำงานของ Observable และ Observer

RxJava และ ReactiveX คืออะไร

RxJava คือการนำข้อกำหนด ReactiveX มาใช้ใน Java ซึ่งเป็นไลบรารีสำหรับการเขียนโปรแกรมแบบอะซิงโครนัสโดยใช้สตรีมที่สังเกตได้ (Observable) RxJava 2 เปิดตัวในปี 2016 พร้อมรองรับ Reactive Streams (Flowable) และการแยกเป็น rx.Observable และ io.reactivex.Observable RxJava 3 (2019) เป็นเวอร์ชันหลักปัจจุบันที่เข้ากันได้ย้อนหลังกับ RxJava 2

แนวคิดหลักของ RxJava คือทุกสิ่งเป็นสตรีม: สตรีมข้อมูล สตรีมเหตุการณ์ สตรีมสถานะ การดำเนินการแบบอะซิงโครนัสใดๆ สามารถแสดงเป็น Observable ที่ส่งข้อมูล ข้อผิดพลาด หรือสัญญาณเสร็จสมบูรณ์ Observer สมัครสมาชิก Observable และได้รับการแจ้งเตือนแบบเรียลไทม์

ตามข้อมูลของ Badoo (2024) ก่อนการเปลี่ยนไปใช้ coroutines 76% ของแอป Android ใน 200 อันดับแรกของ Google Play ใช้ RxJava สำหรับการดำเนินการแบบอะซิงโครนัส ส่วนแบ่งกำลังลดลงเนื่องมาจาก coroutines แต่ RxJava ยังคงอยู่ในโค้ดการผลิตของแอปนับพันและถือเป็นเทคโนโลยีที่เติบโตเต็มที่และผ่านการทดสอบแล้ว ReactiveX เป็นข้อกำหนดข้ามแพลตฟอร์มที่นำไปใช้กับ JavaScript (RxJS), .NET (Rx.NET), Swift (RxSwift) และภาษาอื่นๆ ด้วย

รูปแบบ Observer ใน RxJava

ReactiveX ขยายรูปแบบ Observer แบบคลาสสิกด้วยสองกลไก: การเชื่อมโยงโอเปอเรเตอร์ และ การจัดการเธรดตาม Scheduler Observable จะเริ่มส่งข้อมูลเมื่อมี Observer สมัครสมาชิกเท่านั้น (การประเมินแบบขี้เกียจ) ซึ่งช่วยให้สร้างไปป์ไลน์ข้อมูลที่ทำงานเฉพาะเมื่อมีการสมัครสมาชิก

ประเภทของ Observable: Observable, Flowable, Single, Maybe, Completable

Observable — ประเภทพื้นฐานที่ส่ง 0..N องค์ประกอบด้วย onError หรือ onComplete เหมาะสำหรับสตรีมข้อมูลไม่จำกัด — ตัวอย่างเช่น เหตุการณ์คลิกหรือการอัปเดตตำแหน่งทางภูมิศาสตร์ Observable ไม่รองรับ backpressure

Flowable — เวอร์ชัน Reactive Streams ของ Observable ที่รองรับ backpressure ใช้เมื่อแหล่งข้อมูลสามารถสร้างองค์ประกอบได้เร็วกว่าที่ Observer จะประมวลผล Flowable รองรับกลยุทธ์ BACKPRESSURE_BUFFER, DROP, LATEST และ ERROR

ประเภทองค์ประกอบBackpressureการใช้งาน
Observable0..Nไม่เหตุการณ์ UI, สตรีมขนาดเล็ก
Flowable0..Nใช่ข้อมูลขนาดใหญ่, เวลาจริง
Single1 (onSuccess/onError)การตอบสนองเดี่ยว (เครือข่าย)
Maybe0..1ค่าไม่บังคับ (แคช)
Completable0 (onComplete/onError)การดำเนินการไม่มีข้อมูล (เขียน)

Single, Maybe และ Completable

Single ส่งองค์ประกอบหนึ่งรายการหรือข้อผิดพลาด — เหมาะสำหรับคำขอเครือข่าย Maybe ส่ง 0 หรือ 1 องค์ประกอบ เหมาะสำหรับแคชที่ข้อมูลอาจไม่มี Completable ส่งเฉพาะ onComplete หรือ onError โดยไม่มีข้อมูล สะดวกสำหรับการดำเนินการเขียนหรือลบ ประเภทเหล่านี้ทำให้ API ง่ายขึ้นโดยจำกัดข้อตกลงให้เฉพาะกรณี Retrofit (ไคลเอนต์ HTTP ยอดนิยมสำหรับ Android) รองรับทั้งห้าประเภท RxJava โดยตรง ทำให้สามารถเลือกประเภทส่งคืนที่เหมาะสมที่สุดสำหรับแต่ละเอนด์พอยต์โดยไม่ต้องใช้โค้ดเพิ่มเติม

โอเปอเรเตอร์ RxJava: การแปลงและการกรองสตรีม

โอเปอเรเตอร์ คือฟังก์ชันที่แปลง Observable หนึ่งเป็นอีกอันหนึ่ง ลูกโซ่ของโอเปอเรเตอร์อธิบายไปป์ไลน์ข้อมูล: แต่ละโอเปอเรเตอร์รับสตรีมจากอันก่อนหน้า แปลงมัน และส่งต่อไปยังอันถัดไป RxJava มีโอเปอเรเตอร์มากกว่า 200 ตัวแบ่งตามหมวดหมู่

  • map — แปลงแต่ละองค์ประกอบ (Integer → String)
  • flatMap — แปลงองค์ประกอบเป็น Observable และรวมทั้งหมดเป็นสตรีมเดียว
  • filter — ผ่านองค์ประกอบที่ตรงตามเงื่อนไข
  • zip — รวมองค์ประกอบจาก N Observables ตามดัชนี
  • merge — รวมหลาย Observables เป็นอันเดียว รักษาลำดับเวลา
  • debounce — ส่งองค์ประกอบเฉพาะเมื่อผ่านช่วงเวลาที่กำหนดโดยไม่มีการส่งอื่น

flatMap เป็นหนึ่งในโอเปอเรเตอร์ที่ทรงพลังที่สุดของ RxJava ช่วยให้ดำเนินการคำขอแบบอะซิงโครนัสสำหรับแต่ละองค์ประกอบและรวบรวมผลลัพธ์เป็นสตรีมร่วม ตัวอย่างเช่น flatMap ใช้สำหรับโหลดรายละเอียดจากรายการ ID: แต่ละ ID → คำขอเครือข่าย → การรวมผลลัพธ์ ต่างจาก map ที่เพียงแค่แปลงองค์ประกอบ flatMap สามารถส่งหลายองค์ประกอบหรือสลับไปยัง Observable อื่น ทำให้เป็นพื้นฐานสำหรับการสร้างไปป์ไลน์แบบอะซิงโครนัส

การจัดการข้อผิดพลาดด้วยโอเปอเรเตอร์

onErrorResumeNext — สลับไปยัง Observable สำรองเมื่อเกิดข้อผิดพลาด retry — สมัครสมาชิกใหม่ N ครั้งเมื่อเกิดข้อผิดพลาด onErrorReturn — ส่งค่าปริยายแทนข้อผิดพลาด doOnError — ดำเนินการข้างเคียงเมื่อเกิดข้อผิดพลาดโดยไม่เปลี่ยนสตรีม (การบันทึกหรือการวิเคราะห์) การรวมโอเปอเรเตอร์เหล่านี้ช่วยสร้างไปป์ไลน์ที่แข็งแกร่งด้วยกลยุทธ์การจัดการข้อผิดพลาดที่ชัดเจนโดยไม่ต้องใช้ try/catch ด้วยตนเอง

Schedulers: การจัดการเธรดใน RxJava

Schedulers กำหนดว่า Observable และ Observer ทำงานบนเธรดใด subscribeOn ตั้งค่าเธรดสำหรับแหล่งข้อมูล observeOn ตั้งค่าเธรดสำหรับ Observer และโอเปอเรเตอร์ที่ตามมา การแยกนี้เป็นข้อได้เปรียบหลักของ RxJava: แหล่งข้อมูลบนเธรด IO การประมวลผลบน computation UI บนเธรดหลัก

Schedulers หลัก: Schedulers.io() — สำหรับการดำเนินการ I/O (เครือข่าย, ดิสก์), พูลไม่จำกัด Schedulers.computation() — สำหรับการคำนวณ, พูลคงที่ตามจำนวนคอร์ Schedulers.newThread() — เธรดใหม่สำหรับแต่ละงาน AndroidSchedulers.mainThread() — เธรดหลักของ Android (RxAndroid) นอกจากนี้ยังมี Schedulers.trampoline() สำหรับดำเนินงานในเธรดปัจจุบันด้วยคิว FIFO มีประโยชน์สำหรับการทดสอบ

ตามข้อมูลของ Google (2025) การใช้ Schedulers อย่างถูกต้องเป็นส่วนที่ยากที่สุดของ RxJava สำหรับผู้เริ่มต้น ข้อผิดพลาดทั่วไปคือการเรียก subscribeOn หลังจาก observeOn ซึ่งไม่ส่งผลต่อแหล่งข้อมูล subscribeOn ควรเป็นอันแรกในลูกโซ่สำหรับแหล่งข้อมูล observeOn ก่อนการสมัครสมาชิก UI กฎ: subscribeOn ส่งผลต่อต้นน้ำ (แหล่งข้อมูล) เท่านั้น observeOn สลับปลายน้ำ (ผู้สมัครสมาชิกและโอเปอเรเตอร์ทั้งหมดหลังจากนั้น)

ตัวอย่างโค้ด RxJava ใน Android

พิจารณาสามสถานการณ์: คำขอเครือข่ายด้วย Single คำขอแบบขนานด้วย zip และ debounce สำหรับช่องค้นหาด้วย debounce

คำขอเครือข่ายด้วย Single

Single เหมาะสำหรับคำขอ Retrofit: หนึ่งคำขอ — หนึ่งการตอบสนอง สมัครสมาชิกบนเธรดหลักสำหรับการอัปเดต 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); }
    })

คำขอแบบขนานด้วย zip

zip รวมผลลัพธ์ของ Single อิสระสองอันเป็นอันเดียว ดำเนินการแบบขนาน ผลลัพธ์จะเกิดขึ้นหลังจากทั้งสองเสร็จสมบูรณ์

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 สำหรับช่องค้นหา

debounce ไม่สนใจการเปลี่ยนแปลงข้อความที่รวดเร็วและส่งคำขอหลังจากหยุด 400 ms เท่านั้น distinctUntilChanged ยกเลิกคำขอหากข้อความไม่เปลี่ยนแปลง

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 เทียบกับ Kotlin Coroutines: เปรียบเทียบแนวทาง

RxJava และ Kotlin Coroutines แก้ปัญหาเดียวกัน — การเขียนโปรแกรมแบบอะซิงโครนัส — แต่ด้วยแนวทางที่แตกต่างกันโดยพื้นฐาน RxJava สร้างบนรูปแบบ Observer และเป็นแบบพุช: แหล่งข้อมูลส่งข้อมูล Observer ตอบสนอง Coroutines เป็นแบบดึง: โค้ดขอข้อมูลตามลำดับผ่าน await

  • RxJava — เชิงรีแอคทีฟ, สตรีมข้อมูล, >200 โอเปอเรเตอร์, แบบพุช, เส้นทางการเรียนรู้ชัน
  • Coroutines — ตามลำดับ, suspend/await, ~40 ฟังก์ชัน, แบบดึง, ไวยากรณ์ง่าย
  • RxJava — เติบโตเต็มที่ (2016), ระบบนิเวศขนาดใหญ่, แต่เส้นทางการเรียนรู้ชัน
  • Coroutines — ทันสมัย (2018), ตัวเลือกที่ Google ชื่นชอบสำหรับโค้ดใหม่
  • RxJava — backpressure ในตัวผ่าน Flowable, กลยุทธ์การบัฟเฟอร์ที่ผ่านการทดสอบ
  • Coroutines — Flow พร้อม backpressure เพิ่งมีไม่นาน แต่กำลังพัฒนาอย่างแข็งขันโดย JetBrains

ตามข้อมูลของ Google I/O 2024 Kotlin Coroutines เป็นแนวทางที่แนะนำสำหรับโค้ดแบบอะซิงโครนัสใหม่ใน Android RxJava ยังคงรองรับสำหรับโปรเจกต์ที่มีอยู่ Google มีไลบรารีสะพานเชื่อม (kotlinx-coroutines-rx3) สำหรับการย้ายข้อมูลแบบค่อยเป็นค่อยไป AndroidX (LiveData, Room, Paging 3) รองรับทั้งสองแนวทาง ทำให้สามารถใช้ RxJava ในโมดูลเก่าและ coroutines ในโมดูลใหม่โดยไม่มีความขัดแย้งของ dependencies

กลยุทธ์การย้ายจาก RxJava ไปยัง Coroutines

การเปลี่ยนผ่านแบบค่อยเป็นค่อยไป: แต่ละคอมโพเนนต์ใหม่เขียนด้วย coroutines โค้ด RxJava เก่าไม่ถูกแตะต้อง RxJava → coroutines ผ่าน awaitSingle() หรือ awaitFirst() Coroutines → RxJava ผ่าน future() หรือ asFlowable() การย้ายที่สมบูรณ์ใช้เวลา 6–18 เดือนสำหรับโปรเจกต์ขนาดใหญ่

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

Observable แตกต่างจาก Flowable อย่างไร?

Observable ไม่รองรับ backpressure — หากแหล่งข้อมูลสร้างข้อมูลเร็วกว่าที่ตัวจัดการจะประมวลผล จะเกิด MissingBackpressureException Flowable รองรับ backpressure ของ Reactive Streams ด้วยกลยุทธ์การบัฟเฟอร์ที่ปรับแต่งได้

subscribeOn และ observeOn คืออะไร?

subscribeOn ตั้งค่า Scheduler สำหรับการดำเนินการ Observable ต้นทาง observeOn ตั้งค่า Scheduler สำหรับ Observer และโอเปอเรเตอร์ที่ตามมาทั้งหมดในลูกโซ่ subscribeOn ส่งผลต่อต้นน้ำ observeOn ส่งผลต่อปลายน้ำ

ฉันควรเปลี่ยนจาก RxJava เป็น coroutines หรือไม่?

สำหรับโปรเจกต์ใหม่ — ใช่ Google แนะนำ coroutines สำหรับโปรเจกต์ที่มีอยู่ — ย้ายข้อมูลแบบค่อยเป็นค่อยไปผ่าน kotlinx-coroutines-rx3 RxJava ยังคงเสถียรและรองรับสำหรับโค้ดเก่า

วิธีจัดการข้อผิดพลาดใน RxJava?

ผ่านโอเปอเรเตอร์: onErrorReturn (ค่าปริยาย), onErrorResumeNext (Observable สำรอง), retry (ลองใหม่ N ครั้ง) หรือผ่าน Observer.onError() เพื่อแสดงต่อผู้ใช้

CompositeDisposable คืออะไร?

CompositeDisposable คือคอนเทนเนอร์สำหรับจัดการหลายการสมัครสมาชิก เมื่อเรียกใช้ dispose() การสมัครสมาชิกทั้งหมดที่เพิ่มไว้จะถูกยกเลิก ใช้ใน Activity/Fragment เพื่อยกเลิกคำขอทั้งหมดเมื่อหน้าจอถูกทำลาย

สรุป

  • RxJava — ไลบรารีการเขียนโปรแกรมเชิงรีแอคทีฟสำหรับ Java และ Android บนพื้นฐานรูปแบบ Observer
  • Observable/Flowable — แหล่งข้อมูลที่มีและไม่มีการรองรับ backpressure
  • Single, Maybe, Completable — ประเภทเฉพาะสำหรับ 1, 0..1 และ 0 องค์ประกอบ
  • โอเปอเรเตอร์ (map, flatMap, zip, filter) — ลูกโซ่การแปลงด้วยฟังก์ชันมากกว่า 200 ฟังก์ชัน
  • Schedulers — subscribeOn สำหรับแหล่งข้อมูลและ observeOn สำหรับผู้บริโภค
  • RxJava เทียบกับ Coroutines — Google แนะนำ coroutines สำหรับโค้ดใหม่ RxJava สำหรับโค้ดเก่า
  • CompositeDisposable — การจัดการการสมัครสมาชิกที่ปลอดภัยพร้อมการยกเลิกเมื่อหน้าจอถูกทำลาย

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

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

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

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