RxJava: মৌলিক ধারণা, ReactiveX এবং ডেটা স্ট্রিম নিয়ে কাজ করা

লেখক: IT Sectr প্রকাশিত: 2026-03-16 পড়ার সময়: 8 মিনিট

RxJava হল Java এবং Android-এর জন্য একটি রিঅ্যাকটিভ প্রোগ্রামিং লাইব্রেরি যা Observable এবং Observer-এর মাধ্যমে Observer প্যাটার্ন বাস্তবায়ন করে। ReactiveX GitHub, 2026 অনুসারে, RxJava অপারেটর চেইনের মাধ্যমে অ্যাসিঙ্ক্রোনাস ডেটা স্ট্রিম এবং ইভেন্ট হ্যান্ডেল করতে সক্ষম করে। মৌলিক একক হল Observable, যা ট্রান্সফরমেশনের একটি চেইনের মাধ্যমে Observer-কে ডেটা নির্গত করে। RxJava 3 বর্তমান স্থিতিশীল সংস্করণ যা Java 8 lambda, Reactive Streams এবং RxAndroid-এর মাধ্যমে Android ইন্টিগ্রেশন সমর্থন করে।

মূল বিষয়

  • 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 হিসেবে উপস্থাপন করা যেতে পারে যা ডেটা, ত্রুটি বা সম্পূর্ণতা সংকেত নির্গত করে। একটি Observable-এ একজন Observer সাবস্ক্রাইব করে এবং রিয়েল টাইমে নোটিফিকেশন গ্রহণ করে।

Badoo (2024) অনুসারে, coroutines-এ রূপান্তরের আগে, Google Play-র শীর্ষ 200-এর মধ্যে 76% Android অ্যাপ অ্যাসিঙ্ক্রোনাস অপারেশনের জন্য RxJava ব্যবহার করত। এখন এই অংশ coroutines-এর পক্ষে হ্রাস পাচ্ছে, কিন্তু RxJava হাজার হাজার অ্যাপের প্রোডাকশন কোডে রয়ে গেছে এবং একটি পরিণত, পরীক্ষিত প্রযুক্তি হিসেবে বিবেচিত হয়। ReactiveX হল একটি ক্রস-প্ল্যাটফর্ম স্পেসিফিকেশন যা JavaScript (RxJS), .NET (Rx.NET), Swift (RxSwift) এবং অন্যান্য ভাষার জন্যও বাস্তবায়িত হয়েছে।

RxJava-তে Observer প্যাটার্ন

ReactiveX ক্লাসিক Observer প্যাটার্নকে দুটি মেকানিজম দিয়ে প্রসারিত করে: অপারেটর চেইনিং এবং Scheduler-ভিত্তিক থ্রেডিং। Observable ততক্ষণ পর্যন্ত ডেটা নির্গত করা শুরু করে না যতক্ষণ না একজন Observer সাবস্ক্রাইব করে (লেজি ইভালুয়েশন)। এটি একটি ডেটা পাইপলাইন তৈরি করতে দেয় যা শুধুমাত্র সাবস্ক্রিপশন থাকলেই সক্রিয় হয়।

Observable প্রকার: Observable, Flowable, Single, Maybe, Completable

Observable — onError বা onComplete সহ 0..N এলিমেন্ট নির্গতকারী বেস টাইপ। সীমাহীন ডেটা স্ট্রিমের জন্য উপযুক্ত — উদাহরণস্বরূপ, ক্লিক ইভেন্ট বা জিওলোকেশন আপডেট। Observable ব্যাকপ্রেশার সমর্থন করে না।

Flowable — ব্যাকপ্রেশার সমর্থন সহ Observable-এর Reactive Streams সংস্করণ। এটি ব্যবহৃত হয় যখন ডেটা উৎস Observer-এর প্রসেসিং গতির চেয়ে দ্রুত এলিমেন্ট তৈরি করতে পারে। Flowable BACKPRESSURE_BUFFER, DROP, LATEST এবং ERROR কৌশল সমর্থন করে।

প্রকারএলিমেন্টব্যাকপ্রেশারব্যবহার
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 (Android-এর জন্য একটি জনপ্রিয় HTTP ক্লায়েন্ট) সরাসরি পাঁচটি RxJava টাইপ সমর্থন করে, যা অতিরিক্ত কোড ছাড়াই প্রতিটি এন্ডপয়েন্টের জন্য সবচেয়ে উপযুক্ত রিটার্ন টাইপ বেছে নিতে দেয়।

RxJava অপারেটর: স্ট্রিম ট্রান্সফরমেশন এবং ফিল্টারিং

অপারেটর হল ফাংশন যা একটি Observable-কে অন্যটিতে রূপান্তরিত করে। একটি অপারেটর চেইন ডেটা পাইপলাইন বর্ণনা করে: প্রতিটি অপারেটর পূর্ববর্তী থেকে স্ট্রিম নেয়, এটি রূপান্তরিত করে এবং পরবর্তীতে প্রেরণ করে। RxJava-তে বিভাগগুলিতে গ্রুপ করা 200-এর বেশি অপারেটর রয়েছে।

  • map — প্রতিটি এলিমেন্ট রূপান্তরিত করে (Integer → String)
  • flatMap — একটি এলিমেন্টকে Observable-এ রূপান্তরিত করে এবং সবগুলোকে এক স্ট্রিমে মার্জ করে
  • filter — শর্ত পূরণ করে এমন এলিমেন্ট পাস করে
  • zip — N Observable থেকে এলিমেন্ট ইনডেক্স অনুযায়ী কম্বাইন করে
  • merge — একাধিক Observable-কে একটিতে মার্জ করে, কালানুক্রমিক ক্রম বজায় রাখে
  • 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-র সবচেয়ে কঠিন অংশ। একটি সাধারণ ভুল হল observeOn-এর পরে subscribeOn কল করা, যা উৎসকে প্রভাবিত করে না। subscribeOn উৎসের জন্য চেইনে প্রথম হওয়া উচিত, observeOn UI সাবস্ক্রিপশনের আগে। নিয়ম: subscribeOn শুধুমাত্র upstream (উৎস) প্রভাবিত করে, observeOn downstream (সাবস্ক্রাইবার এবং এর পরে সমস্ত অপারেটর) সুইচ করে।

Android-এ RxJava কোড উদাহরণ

তিনটি পরিস্থিতি বিবেচনা করুন: Single দিয়ে নেটওয়ার্ক অনুরোধ, zip দিয়ে সমান্তরাল অনুরোধ, এবং 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 দ্রুত টেক্সট পরিবর্তন উপেক্ষা করে এবং শুধুমাত্র 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 প্যাটার্নের উপর নির্মিত এবং push-ভিত্তিক: উৎস ডেটা পাঠায়, Observer প্রতিক্রিয়া জানায়। Coroutines pull-ভিত্তিক: কোড await-এর মাধ্যমে ক্রমিকভাবে ডেটা অনুরোধ করে।

  • RxJava — রিঅ্যাকটিভ, ডেটা স্ট্রিম, >200 অপারেটর, push-ভিত্তিক, কঠিন শেখার বক্ররেখা
  • Coroutines — ক্রমিক, suspend/await, ~40 ফাংশন, pull-ভিত্তিক, সহজ সিনট্যাক্স
  • RxJava — পরিণত (2016), বিশাল ইকোসিস্টেম, কিন্তু কঠিন শেখার বক্ররেখা
  • Coroutines — আধুনিক (2018), নতুন কোডের জন্য Google-এর পছন্দের পদ্ধতি
  • RxJava — Flowable-এর মাধ্যমে বিল্ট-ইন ব্যাকপ্রেশার, ভালোভাবে পরীক্ষিত বাফারিং কৌশল
  • Coroutines — ব্যাকপ্রেশার সহ Flow সম্প্রতিকালের, কিন্তু JetBrains দ্বারা সক্রিয়ভাবে উন্নয়নশীল

Google I/O 2024 অনুসারে, Kotlin Coroutines Android-এ নতুন অ্যাসিঙ্ক্রোনাস কোডের জন্য প্রস্তাবিত পদ্ধতি। RxJava বিদ্যমান প্রকল্পগুলির জন্য সমর্থিত রয়েছে। Google ক্রমিক মাইগ্রেশনের জন্য ব্রিজিং লাইব্রেরি (kotlinx-coroutines-rx3) সরবরাহ করে। AndroidX (LiveData, Room, Paging 3) উভয় পদ্ধতি সমর্থন করে, নির্ভরতা দ্বন্দ্ব ছাড়াই পুরানো মডিউলে RxJava এবং নতুনে coroutines ব্যবহারের অনুমতি দেয়।

RxJava থেকে Coroutines-এ মাইগ্রেশন কৌশল

ক্রমিক রূপান্তর: প্রতিটি নতুন কম্পোনেন্ট coroutines দিয়ে লেখা হয়, পুরানো RxJava কোড অপরিবর্তিত রাখা হয়। RxJava → coroutines awaitSingle() বা awaitFirst()-এর মাধ্যমে। Coroutines → RxJava future() বা asFlowable()-এর মাধ্যমে। বড় প্রকল্পগুলির জন্য সম্পূর্ণ মাইগ্রেশনে 6–18 মাস সময় লাগে।

সচরাচর জিজ্ঞাসিত প্রশ্ন

Observable Flowable থেকে কীভাবে আলাদা?

Observable ব্যাকপ্রেশার সমর্থন করে না — যদি উৎস হ্যান্ডলারের প্রসেসিং গতির চেয়ে দ্রুত ডেটা তৈরি করে, তাহলে MissingBackpressureException ঘটে। Flowable কনফিগারযোগ্য বাফারিং কৌশল সহ Reactive Streams ব্যাকপ্রেশার সমর্থন করে।

subscribeOn এবং observeOn কী?

subscribeOn উৎস Observable এক্সিকিউট করার জন্য Scheduler সেট করে। observeOn চেইনে Observer এবং পরবর্তী সব অপারেটরের জন্য Scheduler সেট করে। subscribeOn upstream প্রভাবিত করে, observeOn downstream প্রভাবিত করে।

আমার কি RxJava থেকে coroutines-এ সুইচ করা উচিত?

নতুন প্রকল্পের জন্য — হ্যাঁ, Google coroutines সুপারিশ করে। বিদ্যমান প্রকল্পের জন্য — kotlinx-coroutines-rx3-এর মাধ্যমে ক্রমিক মাইগ্রেশন। RxJava পুরানো কোডের জন্য স্থিতিশীল এবং সমর্থিত রয়েছে।

RxJava-তে ত্রুটিগুলি কীভাবে পরিচালনা করবেন?

অপারেটরের মাধ্যমে: onErrorReturn (ডিফল্ট মান), onErrorResumeNext (ব্যাকআপ Observable), retry (N বার পুনঃচেষ্টা)। অথবা ব্যবহারকারীকে দেখানোর জন্য Observer.onError()-এর মাধ্যমে।

CompositeDisposable কী?

CompositeDisposable একাধিক সাবস্ক্রিপশন পরিচালনার জন্য একটি ধারক। যখন dispose() কল করা হয়, সমস্ত যোগ করা সাবস্ক্রিপশন বাতিল হয়ে যায়। এটি Activity/Fragment-এ স্ক্রিন ধ্বংস হলে সমস্ত অনুরোধ বাতিল করতে ব্যবহৃত হয়।

সারসংক্ষেপ

  • RxJava — Observer প্যাটার্নের উপর ভিত্তি করে Java এবং Android-এর জন্য রিঅ্যাকটিভ প্রোগ্রামিং লাইব্রেরি
  • Observable/Flowable — ব্যাকপ্রেশার সমর্থন সহ এবং ছাড়া ডেটা উৎস
  • Single, Maybe, Completable — 1, 0..1 এবং 0 এলিমেন্টের জন্য বিশেষায়িত টাইপ
  • অপারেটর (map, flatMap, zip, filter) — 200-এর বেশি ফাংশন সহ রূপান্তর চেইন
  • Schedulers — উৎসের জন্য subscribeOn এবং ভোক্তার জন্য observeOn
  • RxJava বনাম Coroutines — Google নতুন কোডের জন্য coroutines সুপারিশ করে, RxJava পুরানোর জন্য
  • CompositeDisposable — স্ক্রিন ধ্বংস হলে বাতিলকরণ সহ নিরাপদ সাবস্ক্রিপশন ব্যবস্থাপনা

আমরা একটি মোবাইল অ্যাপ্লিকেশন টার্নকি তৈরি করব

IT Sectr 2017 সাল থেকে স্টার্টআপ এবং ব্যবসার জন্য iOS এবং Android অ্যাপ্লিকেশন তৈরি করে। আমরা আপনাকে পরামর্শ দেব এবং সেরা সমাধান প্রস্তাব করব।

প্রকল্প নিয়ে আলোচনা করুন

আরও পড়ুন