2. یک مثال مقدماتی
اولین برخوردهای من با RxJava از طریق دورهها و آموزشهایی بود که در اینترنت پیدا کردم. علاوه بر اینکه نظریهای که استفاده میکرد، مفاهیمی را به کار میبرد که با آنها آشنا نبودم و درکشان برایم دشوار بود، واقعاً نمیتوانستم ببینم در زندگی واقعی چه کاربردی میتواند داشته باشد. بنابراین با ارائه یک مثال (امیدوارم سادهای باشد) شروع میکنیم که در آن استفاده از RxJava واقعاً فرآیند نوشتن کد را ساده میکند، و از آنجا سعی میکنیم ویژگیهای کلیدی این کتابخانه را شناسایی کنیم.
کتابخانه RxJava بر اساس مفهوم زیر است: یک جریان از عناصر از نوع T، Observable<T>، توسط یک یا چند مشترک (مشترکها، ناظران، مصرفکنندگان)، Subscriber<T>، مشاهده میشود. کتابخانه RxJava این امکان را فراهم میکند که جریان Observable<T> در یک نخ T1 و ناظر Subscriber<T> آن در یک نخ T2 بدون نیاز به نگرانی توسعهدهنده در مورد مدیریت چرخه عمر این نخها و مسائل دشوار طبیعی، مانند اشتراکگذاری دادهها بین نخها و همگامسازی آنها برای اجرای یک وظیفه کلی، اجرا شود. بنابراین این امر برنامهنویسی ناهمزمان را آسانتر میکند.نیازی به نگرانی در مورد مدیریت چرخه عمر این نخها یا رسیدگی به مسائل ذاتاً دشوار، مانند اشتراکگذاری دادهها بین نخها و همگامسازی آنها برای اجرای یک وظیفه کلی، ندارد. بنابراین، این امر برنامهنویسی ناهمزمان را تسهیل میکند.
یک جریان Observable<T> عناصر از نوع T را تولید میکند که میتوان آنها را در زمان تولیدشان مشاهده کرد. اگر ناظر و قابل مشاهده (اصطلاحی کلی برای اشاره به نوع Observable<T>) در یک نخ باشند، آنگاه قابل مشاهده تنها پس از مصرف عنصر i توسط ناظر، میتواند عنصر (i+1) را تولید کند. موارد اندکی وجود دارد که این معماری کاربرد دارد. اگر ناظر و قابل مشاهده در یک نخ نباشند، آنگاه قابل مشاهده و ناظرش به طور مستقل عمل میکنند: قابل مشاهده با سرعت خود عناصر را منتشر میکند و ناظر با سرعت خود آنها را مصرف میکند. ارزش کتابخانه در همینجاست. تا اینجا، ما فقط در مورد یک ناظر صحبت کردهایم. در واقع، یک قابل مشاهده میتواند هر تعداد ناظری داشته باشد.
2.1. معماری نمونهٔ برنامه
اپلیکیشن مثال معماری زیر را دارد:

- در [1]، یک لایه خدماتی لیستهای اعداد تصادفی را تولید میکند. این لایه در همان تریدِ متد [swing] که از آن استفاده میکند، اجرا میشود. بنابراین اعداد خود را به صورت همزمان تولید میکند؛
- در [2]، یک لایه تطبیق نازک که با RxJava پیادهسازی شده است، این امکان را فراهم میکند که لایه [swing] با یک پیادهسازی ناهمزمان از همان سرویس ارائه شود: این میتواند در یک نخ (thread) متفاوت از نخ متد [swing] که از آن استفاده میکند، اجرا شود؛
- فراخوانی [4] همزمان است، در حالی که فراخوانی [5-6] غیرهمزمان است؛
آنچه میخواهیم در اینجا نشان دهیم این است که کتابخانه Rx تبدیل یک رابط همگام به یک رابط غیرهمگام را آسان میکند. چرا این موضوع مفید است؟ رویدادها در یک رابط Swing در یک نخ که معمولاً به آن حلقه رویداد (event loop) گفته میشود، پردازش میشوند. رویدادها صفبندی شده و یکی پس از دیگری پردازش میشوند. رویداد Ei+1 تا زمانی که رویداد قبلی، Ei، به طور کامل پردازش نشده باشد، قابل پردازش نیست. بنابراین مهم است که پردازش رویداد در کمترین زمان ممکن انجام شود تا رابط کاربری گرافیکی همچنان پاسخگو باقی بماند. گاهی اوقات، پردازش یک رویداد میتواند زمان زیادی ببرد. این زمانی اتفاق میافتد که پردازش شامل دسترسی به شبکه باشد. اگر نمیخواهیم رابط کاربری گرافیکی به گونهای که برای کاربر غیرقابل قبول باشد، از کار بیفتد، این عملیات شبکهای باید در نخهای مجزا از حلقه رویداد انجام شوند تا حلقه رویداد آزاد شود. این ما را به حوزه برنامهنویسی همزمان (که در آن چندین نخ به طور موازی اجرا میشوند) میبرد، حوزهای که به حق دشوار محسوب میشود. کتابخانه Rx یک راهحل ساده و شیک برای این مشکل ارائه میدهد.
برای شبیهسازی عملیات زمانبر، سرویس در این مثال اعداد تصادفی خود را پس از یک تأخیر مشخص برمیگرداند، تا بتوانیم رفتار رابط کاربری گرافیکی را مشاهده کنیم.
2.2. L'exécutable
فایل اجرایی برنامه نمونه را میتوان در پوشه [dvp/executables] در بخش مثالها یافت:
![]() | ![]() |
روشهای مختلفی برای اجرای آرشیو [swing-01] وجود دارد که بسته به پیکربندی کامپیوتری که برای اجرای آن استفاده میشود، متفاوت است. برای مثال، میتوانید فرآیند توصیفشده در [1-3] را دنبال کنید. این کار رابط کاربری گرافیکی زیر را نمایش میدهد:
![]() |
- این رابط کاربری دو زبانه نمایش میدهد: [1-2]، یکی با برچسب [Request] برای درخواست به سرویس تولید عدد تصادفی، و دیگری با برچسب [Response] برای نمایش اعداد دریافتی؛
- در [3]، شما مشخص میکنید که چند درخواست میخواهید به سرویس ارسال کنید؛
- در [4]، محدوده تولید عدد مورد نظر ([a,b]) مشخص میشود؛
- در [5]، تعداد مقادیری که توسط سرویس بازگردانده میشود، یک عدد تصادفی در محدوده [minCount, maxCount] خواهد بود که توسط کاربر تعیین شده است؛
- در [6]، قبل از بازگرداندن پاسخ خود، سرویس به مدت delay میلیثانیه منتظر میماند، که در آن delay یک عدد تصادفی در محدوده تعریفشده توسط کاربر [minDelay, maxDelay] است؛
- به طور پیشفرض، لایه [swing] با رابط همگام سرویس ارتباط برقرار خواهد کرد. برای استفاده از لایه غیرهمگام، کاربر باید گزینه [7] را علامت بزند. در این حالت، سرویس تولید در نخهای (threads) جدا از حلقه رویدادهای رابط کاربری (GUI) اجرا خواهد شد. کتابخانه Rx استراتژیهای مختلفی برای تولید این نخها ارائه میدهد. کاربر میتواند استراتژی خود را از طریق [8] انتخاب کند؛
- اعداد با استفاده از دکمه [9] تولید میشوند؛
![]() |
- در [10]، نتایج نمایش داده میشوند. ما ساختار این نتایج را توضیح خواهیم داد؛
- در [11]، تعداد نتایج بهدستآمده؛
- در [12]، زمان اجرای آن به میلیثانیه؛
- در [13]، کاربر این امکان را دارد که اجرای آن را لغو کند؛
هر نتیجه به شکل زیر است:
{"idClient":0,"serviceResponse":{"delay":412,"aleas":[146,115,128,174,159,112,162,127],"executedOn":"RxComputationThreadPool-6"},"observedOn":"AWT-EventQueue-0","requestAt":"02:42:47:708","responseAt":"02:42:52:931"}
- [idClient]: شماره درخواست. توجه داشته باشید که چندین درخواست به سرویس تولید ارسال میشود؛
- [delay]: زمان انتظار به میلیثانیه که سرویس قبل از ارسال نتیجه خود مشاهده کرده است؛
- [aleas]: اعداد تصادفی بازگرداندهشده توسط سرویس؛
- [executedOn]: نام رشتهای که سرویس در آن اجرا شد؛
- [observedOn]: نام تریدی که نتیجه را نمایش داد. با رابط کاربری Swing، این تنها میتواند ترید حلقه رویداد باشد، که در این مورد [AWT-EventQueue-0] است؛
- [requestAt]: زمان درخواست در قالب [heures:minutes:secondes:millisecondes];
- [responseAt]: زمان دریافت نتایج، با همان فرمت؛
اکنون بخشهای کد مرتبط با درک مثال را ارائه میکنیم.
2.3. رابط همگام

لایه سرویس [1] رابط زیر را فراهم میکند:
public interface IService {
// اعداد تصادفی در [a,b]
//n عدد با n عدد تصادفی در بازه [minCount, maxCount] تولید میشوند
// اعداد پس از تأخیری به میلیثانیه تولید میشوند،
// که در آن [delay] یک عدد تصادفی در بازه [minDelay, maxDelay] است
public ServiceResponse getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay);
}
پاسخ [ServiceResponse] به شرح زیر است:
public class ServiceResponse {
//تأخیر سرویس
private int delay;
// اعداد تصادفی
private List<Integer> aleas;
// رشته اجرایی
private String executedOn;
// سازندهها
public ServiceResponse(int delay, List<Integer> aleas) {
executedOn = Thread.currentThread().getName();
this.delay = delay;
this.aleas = aleas;
}
// گیرنده و تنظیمکننده
...
}
پاسخ دارای سه عنصر است:
- خط ۶: اعداد تصادفی تولید شده؛
- خط ۴: زمان انتظار مشاهدهشده توسط سرویس قبل از بازگرداندن نتیجه آن؛
- خط ۸: نخ اجرای سرویس؛
2.4. فراخوانی همگام

اکنون به تفصیل بررسی میکنیم فراخوانی همگام [4] که توسط لایه [swing] به سرویس [1] انجام شده است:
private void doGenerateWithService() {
// آغاز انتظار
beginWaiting();
try {
for (int i = 0; i < nbRequests; i++) {
UiResponse uiResponse = new UiResponse();
uiResponse.setIdClient(i);
uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
uiResponse.setResponseAt();
model.add(0, jsonMapper.writeValueAsString(uiResponse));
jLabelNbReponses.setText(String.valueOf(Integer.parseInt(jLabelNbReponses.getText()) + 1));
}
} catch (JsonProcessingException | RuntimeException e) {
System.out.println(e);
}
// پایان انتظار
endWaiting();
}
- خطوط ۵–۱۲: حلقه اجرای درخواستهای [nbRequests] انجامشده توسط کاربر؛
- خط ۸: [service] پیادهسازی رابط همگام [IService] است که در بخش ۲.۳ توصیف شده است؛
- خط ۱۰: [model] قالبی است که توسط مؤلفه JList از برگه [Response] نمایش داده میشود. عناصر این قالب رشتههای jSON از عناصر نوع [UiResponse] زیر هستند:
public class UiResponse {
//شناسهٔ مشتری
private int idClient;
// پاسخ سرویس
private ServiceResponse serviceResponse;
// نام نخ مشاهده
private String observedOn;
//زمان درخواست
private String requestAt;
//زمان پاسخ
private String responseAt;
// سازندهها
public UiResponse() {
observedOn = Thread.currentThread().getName();
requestAt = getTimeStamp();
}
// متدهای خصوصی
private String getTimeStamp() {
return new SimpleDateFormat("hh:mm:ss:SSS").format(Calendar.getInstance().getTime());
}
// گیرندهها و تنظیمکنندهها
...
}
- خط ۶: پاسخ از سرویس تولید عدد؛
- خط ۴: شماره درخواستی که به آن پاسخ داده میشود؛
- خط ۸: رشتهای که این پاسخ را نمایش میدهد. همانطور که گفته شد، این همیشه رشته حلقه رویداد خواهد بود؛
- خطوط ۱۰ و ۱۲: زمان درخواست و زمان پاسخ؛
2.5. آزمایش فراخوانیهای همگام
ما پیکربندی زیر را اجرا میکنیم:
![]() |
در برگه [Response] نتایج زیر را به دست میآوریم:
![]() |
- در [1-2]، ما واقعاً ۱۰ پاسخ درخواستی را دریافت کردهایم. آنها در اولین موقعیت و به ترتیبی که رسیدند، درج شدهاند. میتوانیم ببینیم که آنها به ترتیب درخواستها دریافت شدهاند؛
- همه آنها در نخ چرخه رویداد [AWT-EventQueue-0] اجرا و نمایش داده شدند. بنابراین، پرسوجوها یکی پس از دیگری در این نخ اجرا شدند. هیچ پرسوجوی همزمان وجود نداشت؛
- آنچه در اینجا قابل مشاهده نیست این است که در حین اجرا، رابط کاربری گرافیکی (GUI) منجمد میشود. به عنوان مثال، هیچ راهی برای دسترسی به زبانه [Response] برای مشاهده پاسخها هنگام دریافت آنها یا برای متوقف کردن اجرا با استفاده از دکمه [Annuler] وجود ندارد. حتی اگر این دکمه روی زبانه [Request] وجود داشت، غیرقابل استفاده بود. این به این دلیل است که در آن صورت دو رویداد وجود میداشت:
- کلیک روی دکمه [Générer]؛
- یک کلیک روی دکمه [Annuler]؛
کلیک روی دکمه [Annuler] تنها پس از اتمام عملی که با کلیک روی دکمه [Générer] آغاز شده بود، پردازش میشود. ما همین حالا دیدیم که این عملیات، نخ حلقه رویداد را در تمام مدت اجرای خود اشغال کرده و بدین ترتیب مانع از رسیدگی به کلیک روی دکمه [Annuler] شد. این معمولاً از آن دسته موقعیتهایی است که Rx میتواند بهبود قابل توجهی ایجاد کند؛
2.6. رابط ناهمزمان و پیادهسازی آن
اکنون به رابط لایه [2] و پیادهسازی آن با استفاده از Rx نگاهی میاندازیم. این موضوع فوراً واضح نخواهد بود. ما صرفاً میخواهیم سادگی کد در این پیادهسازی را برجسته کنیم.
رابط غیرهمزمان به شرح زیر است:
public interface IRxService {
// اعداد تصادفی در [a,b]
// n عدد با n عدد تصادفی در بازه [minCount, maxCount] تولید میشوند
// اعداد پس از تأخیری به میلیثانیه تولید میشوند،
// که در آن [delay] یک عدد تصادفی در بازه [minDelay, maxDelay] است
public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse);
}
تفاوتها از رابط همگام ارائهشده در بخش 2.3 به شرح زیر است:
- کلاس [UiResponse] که در بخش 2.3 ارائه شده است، اکنون یکی از پارامترهای متد [getAleas] (خط 6) است. دلیل این امر آن است که از آنجا که درخواستها اکنون بهصورت موازی اجرا میشوند و سرویس قبل از بازگرداندن نتیجه، به مدت زمانی تصادفی منتظر میماند، پاسخها به ما به ترتیب درخواستها بازگردانده نخواهند شد. بنابراین ما شیء [UiResponse] را که حاوی، در میان اطلاعات دیگر، شماره درخواست است، ارسال میکنیم:
//شناسهٔ مشتری (درخواست)
private int idClient;
// پاسخ سرویس
private ServiceResponse serviceResponse;
// نام نخ مشاهده
private String observedOn;
// زمان درخواست
private String requestAt;
// زمان پاسخ
private String responseAt;
- نوع پاسخ سرویس ناهمزمان از نوع [Observable<UiResponse>] است. نوع [Observable<>] توسط کتابخانه Rx ارائه میشود. نتیجه از نوع [Observable<UiResponse>] نشان میدهد که متد [getAleas] یک جریان از مقادیر از نوع [UiResponse] را فراهم میکند که یکییکی به ناظر خود ارسال میشوند؛
حال بیایید به پیادهسازی این رابط نگاهی بیندازیم:
public class RxService implements IRxService {
// خدمات
private IService service;
// تولیدکننده
public RxService(IService service) {
this.service = service;
}
@Override
public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse) {
return Observable.create(subscriber -> {
try {
uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
subscriber.onNext(uiResponse);
} catch (Exception e) {
subscriber.onError(e);
} finally {
subscriber.onCompleted();
}
});
}
}
- خطوط ۷–۹: یک مرجع به رابط همگام [IService] به سازنده ارائه میشود. این همان چیزی است که تولید اعداد تصادفی را مدیریت خواهد کرد؛
- مشاهدهپذیری که توسط متد [getAleas] بازگردانده میشود، توسط متد استاتیک [Observable.create] ساخته میشود. این متد است که امکان ساخت یک پیادهسازی غیرهمزمان از روی یک پیادهسازی همزمان را فراهم میکند؛
- خط ۱۳: پارامتر متد استاتیک [Observable.create] در اینجا یک تابع لامبدا است که یک نوع [Subscriber] را به عنوان پارامتر خود میگیرد، که این نوع نیز، مجدداً، یک نوع Rx است. یک [Subscriber] شیئی است که به یک جریان از قابلمشاهدهها (observables) مشترک میشود، یعنی جریانی از دادهها که بهصورت ناهمزمان تحویل داده میشوند. سه متد این مشترک در اینجا استفاده شدهاند:
- [Subscriber.onNext] برای ارسال مقداری داده به آن (خط 16)؛
- [Subscriber.onError] برای ارسال یک استثنا به آن (خط 18);
- [Subscriber.onCompleted] برای اطلاعرسانی به مشترک که جریان داده پایان یافته است (خط ۲۰)؛
ممکن است چندین مشترک برای یک مشاهدهپذیر واحد وجود داشته باشد. در اینجا، ما تنها یک مشترک خواهیم داشت که به یک جریان از یک قطعه داده واحد، که در خطوط ۱۵–۱۶ تولید میشود، مشترک میشود. این داده توسط پیادهسازی همزمان سرویس (خط ۱۵) تولید و به مشترک بازگردانده میشود (خط ۱۶).
اگرچه همه این موارد ممکن است هنوز کمی مبهم به نظر برسد، اما نمیتوان از دقت و اختصار فوقالعاده این پیادهسازی غیرهمزمان سرویس چشمپوشی کرد.
2.7. فراخوانی ناهمزمان

اکنون به تفصیل تماس همگام [5] را که توسط لایه [swing] به سرویس [2] ارسال میشود، بررسی خواهیم کرد:
private void doGenerateWithRxService() {
// شروع انتظار
beginWaiting();
// درخواست اعداد تصادفی
Observable<UiResponse> observables = Observable.empty();
for (int i = 0; i < nbRequests; i++) {
UiResponse uiResponse = new UiResponse();
uiResponse.setIdClient(i);
// برنامهریز
int schedulerIndex = jComboBoxSchedulers.getSelectedIndex();
switch (schedulerIndex) {
case 0:
observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.io()));
break;
...
}
}
...
}
- خطوط ۶–۱۰: اجرای درخواستهای [nbRequests] انجامشده توسط کاربر؛
- خطوط ۷–۸: آمادهسازی شیء [UiResponse] مورد نیاز متد [getAleas] سرویس غیرهمزمان (خط ۱۳). این عمدتاً شامل ثبت شماره درخواست [idClient] است؛
- خط ۱۳: متد [getAleas] سرویس ناهمزمان فراخوانی میشود. این متد یک شیء [Observable<UiResponse>] را بازمیگرداند. این فراخوانی هنوز سرویس همزمان را اجرا نمیکند. بیایید به کد سرویس ناهمزمان [getAleas] بازگردیم:
@Override
public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse) {
return Observable.create(subscriber -> {
try {
uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
subscriber.onNext(uiResponse);
} catch (Exception e) {
subscriber.onError(e);
} finally {
subscriber.onCompleted();
}
});
}
کد در خطوط ۴ تا ۱۱، که سرویس همگام را فراخوانی میکند، تنها زمانی اجرا میشود که یک مشترک ثبتنام کند. تا زمانی که هیچ مشترکی وجود نداشته باشد، این کد اجرا نمیشود.
بیایید به کد متد [doGenerateWithRxService] بازگردیم:
- خط ۵: یک observable خالی ایجاد میشود (هیچ چیزی در حال مشاهده نیست)؛
- خط ۱۳: یک observable ایجاد میشود که جریان آن حاصل ادغام جریانهای ناهمزمان [nbRequests] مرتبط با درخواستهای [nbRequests] خواهد بود. این کار با استفاده از متد [Observable.mergeWith] انجام میشود که امکان ادغام دو جریان ناهمزمان را فراهم میکند. در اصطلاح Rx، [mergeWith] یک اپراتور جریان (stream operator) نامیده میشود. ویژگی متمایز این اپراتورها این است که نتیجه عملیات، در اکثر موارد، یک [Observable] دیگر است. در نهایت، پس از خط 17، متغیر [observables] به یک جریان واحد اشاره دارد که شامل پاسخهای غیرهمزمان [nbRequests] تولید شده توسط سرویس غیرهمزمان است؛
- خط ۱۳: عملیات ادغام میتوانست به صورت زیر نوشته شود:
observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse));
اما ما نوشتهایم:
observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.io()));
در اینجا، ما از اپراتور [subscribeOn] بر روی مشاهدهپذیر [rxService.getAleas] استفاده کردهایم. همانطور که اغلب اتفاق میافتد، نتیجه یک بار دیگر یک مشاهدهپذیر است. اپراتور [subscribeOn] به ما امکان میدهد مشخص کنیم که مشاهدهپذیر باید در یک تِرد ارائهشده توسط [Scheduler] اجرا شود. چندین اپراتور [Scheduler] مناسب برای موقعیتهای مختلف وجود دارد. در رابط کاربری گرافیکی، ما تعدادی از آنها را فراهم کردهایم تا بتوانید اثرات هر یک را ببینید:
![]() |
این کد زیر را تولید میکند:
private void doGenerateWithRxService() {
// شروع انتظار
beginWaiting();
// درخواست اعداد تصادفی
Observable<UiResponse> observables = Observable.empty();
for (int i = 0; i < nbRequests; i++) {
UiResponse uiResponse = new UiResponse();
uiResponse.setIdClient(i);
// برنامهریز
int schedulerIndex = jComboBoxSchedulers.getSelectedIndex();
switch (schedulerIndex) {
case 0:
observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.io()));
break;
case 1:
observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.computation()));
break;
case 2:
observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.newThread()));
break;
case 3:
observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.trampoline()));
break;
case 4:
observables = observables.mergeWith(rxService.getAleas(a, b, minCount, maxCount, minDelay, maxDelay, uiResponse).subscribeOn(Schedulers.immediate()));
break;
}
}
...
}
بیایید دوباره به کد در خطوط ۱۲–۱۴ نگاه کنیم. زمانبندیکننده [Schedulers.io()] برای هر مشاهدهپذیر یک نخ جدید اختصاص میدهد. اگر کد را دنبال کنیم:
- خط ۵: یک observable خالی داریم؛
- خط ۱۳، تکرار ۱: observables برابر است با لیست [observable0/thread0] (ناظر observable0 که روی نخ thread0 اجرا شده است)؛
- خط ۱۳، تکرار ۲: observables لیست [observable0/thread0, observable1/thread1] است؛
- و غیره...
در نهایت، پس از خط ۲۸، ما یک مشاهدهپذیر (observable) داریم که از ادغام مشاهدهپذیرهای [nbRequests] حاصل شده و روی نخهای (threads) مختلف [nbRequests] اجرا میشود. همه زمانبندها (schedulers) به این شکل کار نمیکنند، همانطور که در طول آزمایشها خواهیم دید.
بیایید بررسی کد فراخوانی سرویس ناهمزمان را ادامه دهیم:
private void doGenerateWithRxService() {
// شروع انتظار
beginWaiting();
// درخواست اعداد تصادفی
Observable<UiResponse> observables = Observable.empty();
for (int i = 0; i < nbRequests; i++) {
...
}
// ناظر
observables = observables.observeOn(SwingScheduler.getInstance());
// اجرای این مشاهدات
subscriptions.add(observables.subscribe(uiResponse -> {
updateUi(uiResponse);
} , th -> {
System.out.println(th);
doCancel();
} , this::doCancel));
}
- دیدیم که وقتی به خط ۱۰ میرسیم، یک مشاهدهپذیر واحد وجود دارد، که ادغام مشاهدهپذیرهای [nbRequests] است و بسته به زمانبندیکنندهای که کاربر انتخاب کرده، ممکن است روی نخهای [nbRequests] مختلف اجرا شود یا نشود؛
- خط ۱۰: اپراتور [observeOn] به ما امکان میدهد تا مشخص کنیم که میخواهیم دادهها را از کدام تِردِ ناظر بازیابی کنیم — در این مورد، اشیاء [nbRequests] از نوع [UiResponse]. در یک رابط Swing، هیچ انتخابی وجود ندارد. هرگونه بهروزرسانی رابط باید در نخ حلقه رویداد انجام شود. در اینجا، دادهها از مشاهدهپذیر در یک کامپوننت Swing از نوع JList نمایش داده خواهد شد. تایپاسکریپت [SwingScheduler.getInstance()] نمایانگر نخ حلقه رویداد است. کلاس [SwingScheduler] از کتابخانه RxJava منشأ نمیگیرد، بلکه از کتابخانه مشتق RxSwing منشأ میگیرد؛
- وقتی به خط ۱۲ میرسیم، سرویس همگام هنوز فراخوانی نشده است زیرا ناظر در خط ۱۰ هنوز مشترکی ندارد. خطوط ۱۲–۱۷ با استفاده از اپراتور [subscribe] یک مشترک برای آن فراهم میکنند. پارامترهای این اپراتور در این مورد سه تابع لامبدا هستند:
- اولین مورد، [uiResponse -> {updateUi(uiResponse);}]، یکی از اشیاء [UiResponse] تولیدشده توسط مشاهدهپذیر را بهعنوان پارامتر میپذیرد. به یاد داشته باشید که در اینجا، ما اشیاء [nbRequests] از این نوع را خواهیم داشت. متد مرتبط، در این مورد updateUi، باید از این نتیجه استفاده کند؛
- دومین [th -> {System.out.println(th);doCancel();}] یک نوع [Throwable] را بهعنوان پارامتر میپذیرد؛ در این مورد، استثنایی که هنگام اجرای مشاهدهپذیر رخ داده است. متد مربوطه باید از این اطلاعات استفاده کند. در اینجا، این مقدار بر روی کنسول نمایش داده میشود (خط ۱۵) و اجرای برنامه لغو میشود، که این امر برخی از عناصر رابط کاربری گرافیکی را بهروزرسانی خواهد کرد؛
- سومین [this::doCancel] زمانی فراخوانی میشود که مشاهدهپذیر سیگنال دهد که دیگر دادهای برای ارسال ندارد. در اینجا، مشاهدهپذیر، اتحاد مشاهدهپذیرهای [nbRequests] است. مشاهدهپذیر حاصل زمانی که همه مشاهدهپذیرهای تشکیلدهنده آن خودشان سیگنال پایان کار را بدهند، نشان میدهد که کارش تمام شده است. بنابراین، زمانی که این سومین تابع لامبدا اجرا میشود، تمام دادهها دریافت شده است. متد محلی [doCancel] رابط کاربری گرافیکی را برای نشان دادن پایان اجرای برنامه بهروزرسانی میکند؛
متغیر [subscriptions] به صورت زیر تعریف شده است:
// اشتراکها در مشاهدات
protected List<Subscription> subscriptions = new ArrayList<Subscription>();
نوع [Subscription] یک اشتراک را نشان میدهد، یعنی پیوند بین مشترک [Subscriber] و آنچه او مشاهده میکند [Observable]. ما در اینجا از فهرستی از اشتراکها استفاده کردهایم، اگرچه در این مثال تنها یک مورد وجود دارد. متد محلی [doCancel]، که زمانی اجرا میشود که سیگنال قابل مشاهده اعلام کند دیگر دادهای برای ارسال ندارد، به شرح زیر است:
@Override
protected void doCancel() {
// پایان انتظار
endWaiting();
// در مورد اشتراکها
if (jCheckBoxRxSwing.isSelected() && subscriptions != null) {
subscriptions.forEach(Subscription::unsubscribe);
}
}
- خط ۷ همه مشترکین را از مشاهدهپذیر لغو اشتراک میکند؛
از این توضیح مختصر، نکات کلیدی زیر قابل استخراج است:
- نوع [Observable] نشاندهنده یک جریان از مقادیر است که یکییکی به مشترکین یا ناظران ارسال میشوند؛
- نوع [Subscriber] به یک مشترک از نوع [Observable] اشاره دارد؛
- نوع [Subscription] نشاندهنده یک اشتراک است، یعنی پیوند بین یک [Subscriber] و یک [Observable]؛
- نوع [Observable] اپراتورهایی از نوع [mergeWith, empty, subscribeOn, observeOn, ...] را میپذیرد که بیشتر آنها قابلمشاهده تولید میکنند. این اپراتورها برای پیکربندی قابلمشاهده قبل از اجرای آن استفاده میشوند:
- آنچه باید مشاهده شود؛
- رشتهای که مشاهدهپذیر روی آن اجرا میشود؛
- رشتهای که مشترک در آن دادهها را از مشاهدهپذیر دریافت میکند؛
- دو نوع مشاهدهپذیر وجود دارد: [froid / cold] و [chaud / hot]. یک مشاهدهپذیر سرد هر بار که یک مشترک جدید اضافه میشود، به طور کامل اجرا میشود. اگر هر اجرا دادههای یکسانی تولید کند، هر مشترک جدید همان دادههای مشترک قبلی را دریافت میکند. یک مشاهدهپذیر داغ معمولاً دادهها را بهطور مداوم تولید میکند. وقتی یک مشترک، مشترک میشود، دادههایی را که از زمان مشترک شدنش منتشر شدهاند دریافت میکند. او هیچ دادهای را که ممکن است قبلاً منتشر شده باشد، دریافت نمیکند. در مثال ما، قابلیت مشاهدهپذیری سرد است: هر بار که مشترک جدیدی اضافه میشود، به طور کامل مجدداً اجرا میشود. در مثال ما در واقع چه چیزی اجرا میشود؟ برای یافتن پاسخ، باید به تعریف قابلیت مشاهدهپذیری مشاهدهشده بازگردیم:
@Override
public Observable<UiResponse> getAleas(int a, int b, int minCount, int maxCount, int minDelay, int maxDelay, UiResponse uiResponse) {
return Observable.create(subscriber -> {
try {
uiResponse.setServiceResponse(service.getAleas(a, b, minCount, maxCount, minDelay, maxDelay));
subscriber.onNext(uiResponse);
} catch (Exception e) {
subscriber.onError(e);
} finally {
subscriber.onCompleted();
}
});
}
برای هر مشترک جدید، تابع لامبدا—که پارامتری از متد [Observable.create] (خط ۳) است—مجدداً اجرا میشود. بنابراین، خطوط ۴ تا ۱۱ برای هر مشترک جدید در [subscriber] اجرا میشوند؛
2.8. آزمایش فراخوانیهای ناهمزمان
ما با نمایش تأثیر برنامهریزهای مختلف موجود شروع میکنیم. برای این کار، از پارامترهای زیر استفاده میکنیم:
![]() |
ما مقدار [1-2] را روی مقادیر کوچک تنظیم میکنیم تا حتی اگر درخواستها روی همان نخ اجرا شوند، مجبور نباشیم مدت زیادی منتظر بمانیم.
2.8.1. با برنامهریز [Schedulers.io]
![]() |
نکات زیر قابل توجه است:
- پاسخها به ترتیبی متفاوت از درخواستها دریافت میشوند (به idClient مراجعه کنید)؛
- هر درخواست در یک نخ (thread) مجزا اجرا شد؛
- این بار رابط کاربری گرافیکی دیگر منجمد نمیشود:
- میتوانید بین تبها جابجا شوید؛
- میتوانید دادههای در حال ورود را مشاهده کنید؛
- به دلیل سرعت بالای اجرا، فرصتی برای دیدن دکمه [Annuler] وجود ندارد. ما این مورد را در یک تست دیگر برجسته خواهیم کرد؛
2.8.2. با برنامهریز [Schedulers.computation]
![]() |
نکات زیر قابل توجه است:
- پاسخها به ترتیبی متفاوت از درخواستها دریافت میشوند (به idClient مراجعه کنید)؛
- درخواستها در ۸ نخ اجرا شدند؛
- رشته شماره ۳ برای درخواستهای ۸ و ۰ استفاده شد؛
- رشته شماره ۴ برای پرسوجوهای ۹ و ۱ استفاده شد؛
- هر یک از پرسوجوهای دیگر از یک نخ (thread) متفاوت استفاده کردند؛
برنامهریز [Schedulers.computation] به تعداد هستههای ماشین مورد استفاده، نخ ایجاد میکند. این اطلاعات از طریق عبارت [Runtime.getRuntime().availableProcessors()] به دست میآید.
2.8.3. با برنامهریز [Schedulers.newThread]
![]() |
این به روشی مشابه زمانبندیکننده [Schedulers.io] عمل میکند.
2.8.4. با برنامهریزهای [Schedulers.trampoline, Schedulers.immediate]
![]() |
رفتار همگام است. تمام درخواستها روی نخ حلقه رویداد اجرا میشوند. این نتیجه نباید تعمیم داده شود؛ بلکه صرفاً به این معناست که در این مثال خاص، هر دو زمانبندیکننده بهصورت همگام عمل کردهاند.
2.9. موارد مرزی
در این مثال، با برنامهریزهایی کار خواهیم کرد که از عملیات ناهمزمان پشتیبانی میکنند. ابتدا، با استفاده از برنامهریز [Schedulers.computation] که در اینجا روی ۸ نخ اجرا میشود، تعداد درخواستها را به ۱۰۰ افزایش میدهیم. نتیجه زیر را به دست میآوریم:
![]() |
- در [1]، دکمه [Annuler] موجود و قابل استفاده است (عملیات ناهمزمان)؛
اکنون، اجازه دهید اجرای برنامه تا انتها ادامه یابد:
![]() |
از [2] میتوان دید که اجرای ۱۰۰ پرسوجو تقریباً ۴ ثانیه طول کشید (در ۸ نخ).
اکنون، بیایید همین ۱۰۰ درخواست را با استفاده از برنامهریز [Schedulers.newThread] اجرا کنیم که هر درخواست را در یک نخ جداگانه اجرا میکند:
![]() |
در [1]، میبینیم که اجرای ۱۰۰ درخواست (در ۱۰۰ نخ) نیم ثانیه طول کشید. بنابراین این به طور قابل توجهی سریعتر از برنامهریز [Schedulers.computation] است.
اکنون، بیایید ۸۰۰ درخواست را تحت همان شرایط، مجدداً با استفاده از زمانبندیکننده [Schedulers.newThread] اجرا کنیم. نتایج زیر به دست میآید:
![]() |
۸۰۰ درخواست تقریباً در ۱ ثانیه اجرا میشوند.
وقتی این تعداد افزایش مییابد (فراتر از ۲۵۰۰ درخواست روی دستگاه من – که در ۱.۵ ثانیه اجرا شد – این رقم البته به شدت به محیط کاری در زمان اجرا بستگی دارد)، در نهایت با استثنای زیر مواجه میشویم:
![]() |
این نشاندهنده سرریز استک (stack overflow) است. آزمایشها نشان میدهند که رفتار برنامهزمانبندیکننده [Schedulers.newThread] قطعی (deterministic) نیست. ممکن است با استثنای فوق مواجه شوید، آزمایشهای بیشتری را اجرا کنید، سپس به پیکربندیای که باعث ایجاد استثنا شده بود بازگردید و ببینید که دیگر رخ نمیدهد.
2.10. Conclusion
ما مثالی از نحوه استفاده از کتابخانه Rx را نشان دادیم. بیایید آنچه را که آموختهایم خلاصه کنیم:
ما با معماری زیر شروع کردیم:

- در [4]، لایه [swing] فراخوانیهای همزمان را به لایه [service] انجام داد؛
- در [5]، لایه [swing] فراخوانیهای ناهمزمان به لایه [rxService] انجام میداد، که به نوبه خود به صورت همزمان لایه [6] را برای لایه [service] فراخوانی میکرد؛
اولین چیزی که متوجه شدیم این بود که کتابخانه Rx ایجاد رابط غیرهمزمان [rxService] از رابط همزمان [service] را آسان میکند (به بخش 2.4 مراجعه کنید). این یک درس مهم است زیرا به این معنی است که یک برنامه همگام به راحتی میتواند به یک برنامه ناهمزمان تبدیل شود.
در لایه [swing]، دو متد جداگانه نوشته شدند:
- یکی برای انجام فراخوانیهای همزمان به سرویس (بخش 2.4 را ببینید)؛
- دیگری برای انجام فراخوانیهای ناهمزمان به آن (ببینید بخش 2.7)؛
نوشتن فراخوانیهای غیرهمزمان به مراتب پیچیدهتر از نوشتن فراخوانیهای همزمان بود. با این حال، کسانی که تجربه برنامهنویسی همزمان با چندین نخ برای همگامسازی دارند، متوجه خواهند شد که راهحل Rx نوشتن سادهتری دارد و از تمام مشکلات دشوار همگامسازی و ارتباط بین نخها جلوگیری میکند. در حین نوشتن این کد، نکات کلیدی زیر را شناسایی کردیم:
- نوع [Observable] نشاندهنده یک جریان از رویدادها (مقادیر) است که ممکن است (اما الزاماً نه) ناهمزمان باشد و قابل مشاهده است؛
- نوع [Subscriber] نمایانگر یک مشترک برای نوع [Observable] است؛
- نوع [Subscription] نشاندهنده یک اشتراک است، یعنی پیوند بین یک [Subscriber] و یک [Observable]؛
- نوع [Observable] اپراتورهای [mergeWith, empty, subscribeOn, observeOn, ...] را میپذیرد که بیشتر آنها قابلمشاهده هستند. این اپراتورها برای پیکربندی قابلمشاهده قبل از اجرای آن استفاده میشوند:
- آنچه باید مشاهده شود؛
- رشتهای که مشاهدهپذیر روی آن اجرا میشود؛
- رشتهای که مشترک در آن دادهها را از مشاهدهپذیر دریافت میکند؛
- دو نوع مشاهدهپذیر وجود دارد: [froid / cold] و [chaud / hot]. یک مشاهدهپذیر سرد هر بار که یک مشترک جدید اضافه میشود، به طور کامل اجرا میشود. اگر هر اجرا دادههای یکسانی تولید کند، هر مشترک جدید همان دادههای مشترک قبلی را دریافت میکند. یک observable داغ (hot) به طور کلی دادهها را به صورت مداوم تولید میکند. هنگامی که یک مشترک (subscriber) مشترک میشود، دادههایی را که از زمان اشتراک او منتشر شدهاند، دریافت میکند. او هیچ دادهای را که ممکن است قبلاً منتشر شده باشد، دریافت نمیکند. در مثال ما، observable سرد است: این observable برای هر مشترک جدید به طور کامل مجدداً اجرا میشود.
اکنون که مثالی را دیدیم که مزایای کتابخانه Rx را نشان میدهد، آن را با جزئیات بیشتری بررسی خواهیم کرد.
کتابخانه Rx حاوی متدهای متعددی است که در امضاهایشان پارامترهای عمومی دارند. ما یک مرور کلی کوتاه از این امضاها (بخش ۳) ارائه خواهیم داد. پارامترهای این متدها عمدتاً رابطهای تابعی (Java 8) هستند، یعنی رابطهایی با تنها یک متد. بنابراین، پارامترهای واقعی باید نمونههایی از این رابطها باشند. پیش از جاوا ۸، متداول بود که یک رابط را با استفاده از یک کلاس ناشناس پیادهسازی کنند. با جاوا ۸، اگر رابط یک رابط تابعی باشد، پیادهسازی آن با استفاده از یک عبارت لامبدا مختصرتر است. بنابراین، ما این موارد را مورد بحث قرار خواهیم داد (بخش ۴). پس از انجام این کار، کلاس [Stream] (بخش ۵) را معرفی خواهیم کرد که امکان پردازش کلکسیونهای جاوا را با استفاده از توابع لامبدا فراهم میکند. این کلاس از آن جهت جالب توجه است که کلاس [Observable]، که از RxJava مشتق شده است، از آن قرض میگیرد:
- روشهای خاص؛
- همان روش زنجیرهسازی متدها برای پردازش یک مشاهدهپذیر واحد؛
سپس رابطهای تابعی مخصوص کتابخانه RxJava (بند 6) را ارائه خواهیم داد. به بررسی عناصر اصلی کتابخانه Rx [Observable, Subscriber, Subscription, opérateurs] (بند 7) ادامه میدهیم. کلاس [Observable] دهها اپراتور دارد که خود بارها و بارها اِبربارگذاری شدهاند. این موضوع در ابتدا پیچیدگی زیادی ایجاد میکند، زیرا این اپراتورها و ابربارهایشان گاهی تنها در یک جزئیات با هم متفاوت هستند و بدون تجربه، دشوار است که بدانیم کدام اپراتور را باید استفاده کرد. ما تنها تعداد محدودی از اپراتورها را ارائه خواهیم داد و در بیشتر موارد، اورلودهای آنها را نادیده میگیریم.
تمام بخش قبلی با استفاده از کتابخانه RxJava در برنامههای کنسول ساده پوشش داده خواهد شد. پس از اینکه با کتابخانه RxJava آشنا شدیم، از آن در دو نوع برنامه گرافیکی استفاده خواهیم کرد:
- در بخش ۸، مثال برنامه Swing را مجدداً بررسی خواهیم کرد تا آن را با جزئیات بیشتری مورد مطالعه قرار دهیم. سپس از کتابخانه RxSwing استفاده خواهیم کرد؛
- در بخش ۹، یک اپلیکیشن اندروید با استفاده از کتابخانه RxAndroid ایجاد خواهیم کرد؛
پس از انجام همه این مراحل، خواننده ابزار لازم را برای ایستادن روی پای خود خواهد داشت. احتمالاً مدتی طول میکشد تا بتواند کتابخانه Rx را بهصورت شهودی بهکار گیرد. من این کتابخانه را بهویژه جالب یافتم. با این حال، درک آن برایم پیچیده بود و منحنی یادگیری آن تند بود. امیدوارم این سند به کوتاهتر شدن آن منحنی یادگیری برای خواننده کمک کند. به نظر من این تلاش کاملاً ارزشمند است.
















