当前位置: 首页 > news >正文

Android RxJava 实战入门:解决异步、线程切换与生命周期绑定三大痛点

1. 这不是又一个“RxJava 概念堆砌”教程——它解决的是 Android 开发者真正卡住的三个具体问题

你点开这个标题,大概率是因为:在 Android Studio 里写了个网络请求,结果 UI 线程被阻塞、页面卡死;或者用Handler+Runnable嵌套了三层回调,改个逻辑要重读五分钟代码;又或者刚学完LiveData,发现它处理“搜索框实时过滤+防抖+取消上一次请求”这种场景时,写起来依然绕得像解九连环。这些不是抽象的“响应式编程思想”,而是每天早上九点打开 IDE 就扑面而来的具体痛点。RxJava 不是银弹,但它确实是目前 Android 生态里,唯一能把异步、线程切换、事件流组合、错误传播这四件事打包成一套可预测、可复用、可测试的语法糖的工具。我带过 7 个 Android 团队,新成员上手最快的不是看官方文档,而是直接跑通一个带“搜索防抖+加载状态+错误重试”的完整链路——因为所有概念都长在真实业务里,而不是飘在Observable.just()的示例里。这篇教程不讲flatMapconcatMap的哲学区别,只告诉你:当用户在搜索框里每秒敲 3 个字,你如何用 4 行代码让后台只发起 1 次有效请求,并且在请求失败时自动弹出带“重试”按钮的 Toast;当你需要把SharedPreferences的变更、数据库插入、网络同步三件事串成“先存本地再发服务器最后刷新 UI”的原子操作,如何避免写 12 个if-else判断每个环节的成功与否。关键词AndroidRxJava教程入门—— 它们指向的不是理论考试,而是你明天就要提交的 PR 里那几行能跑通、能维护、能加单元测试的代码。

2. 为什么 RxJava 在 Android 上不可替代?——从线程模型和生命周期绑定说起

2.1 Android 的线程困境:不是“能不能切”,而是“切完怎么收场”

Android 的主线程(UI 线程)就像一家餐厅的前台服务员:所有用户点击、滚动、动画都得排队等他响应。但如果你在主线程里直接调用OkHttpClient.newCall().execute(),相当于让服务员亲自去后厨炒菜——前台瞬间瘫痪,用户看到的就是白屏或 ANR 对话框。传统解法是AsyncTask(已废弃)、HandlerThreadExecutorService,但它们只解决“把活扔给后台干”,没解决“干完怎么把结果安全送回前台”以及“用户按了返回键,这个后台任务要不要停”。RxJava 的核心价值,恰恰卡在这个缝隙里:它把任务调度(Scheduler)和任务本身(Observable)彻底解耦。你看这段代码:

apiService.search(query) .subscribeOn(Schedulers.io()) // 指定在 IO 线程执行网络请求 .observeOn(AndroidSchedulers.mainThread()) // 指定在主线程接收结果 .subscribe( result -> updateUi(result), // 成功时在主线程更新 UI error -> showError(error) // 错误时也在主线程弹 Toast );

subscribeOnobserveOn不是魔法,它们背后是 RxJava 内置的线程池管理器。Schedulers.io()维护一个可扩容的线程池(默认最大 64 个线程),专为磁盘读写、网络请求这类阻塞操作设计;AndroidSchedulers.mainThread()则利用Handler机制,把回调安全地 post 到主线程消息队列。关键在于:你不需要手动创建HandlerLooperMessageQueue,更不用在onDestroy()里反复removeCallbacks()。RxJava 的Disposable接口(比如CompositeDisposable)会自动帮你清理未完成的任务——这直接对应 Android 的Activity/Fragment生命周期。我见过太多项目,因为忘记在onPause()里取消HandlerpostDelayed,导致 Activity 销毁后还在尝试更新已不存在的 View,最终NullPointerException爆满日志。

2.2 生命周期绑定:为什么CompositeDisposableWeakReference更可靠?

新手常犯的错误是:用WeakReference<Activity>包裹回调,在onDestroy()里清空引用。但WeakReference只解决“内存泄漏”,不解决“逻辑错误”。举个例子:用户在搜索页输入“手机”,RxJava 发起请求;用户立刻按返回键离开页面,Activity被销毁;但网络请求 2 秒后才返回,WeakReference.get()返回 null,updateUi()不执行——看起来没问题?错。如果这个请求还触发了SharedPreferences存储或数据库写入,这些副作用依然会发生,只是 UI 没更新而已。而CompositeDisposable的设计哲学是:“任务要么全部完成,要么全部取消”。它内部维护一个List<Disposable>,调用clear()时,会逐个调用每个Disposable.dispose()方法。对于网络请求,dispose()会直接调用OkHttp Call.cancel();对于定时任务,会移除ScheduledExecutorService中的Future。这意味着:副作用(如数据库写入)根本不会发生,而不是发生了但找不到 UI 更新。我在某电商 App 的商品详情页实测过:用CompositeDisposable管理图片加载和价格查询两个 Observable,用户快速滑动列表时,99.8% 的图片加载请求被及时取消,内存占用比用WeakReference降低 40%。这不是玄学,是 RxJava 把“取消语义”从应用层下沉到了框架层。

2.3 事件流 vs 单次回调:为什么flatMap比嵌套Callback更易维护?

传统回调地狱长这样:

// 获取用户信息 api.getUser(userId, new Callback<User>() { @Override public void onSuccess(User user) { // 根据用户地区获取天气 api.getWeather(user.region, new Callback<Weather>() { @Override public void onSuccess(Weather weather) { // 根据天气推荐商品 api.getRecommendations(weather.season, new Callback<List<Product>>() { @Override public void onSuccess(List<Product> products) { updateRecommendation(products); } // 三个层级的 onError,每个都要写一遍错误处理 }); } }); } });

三层嵌套,5 个@Override,错误处理分散。而 RxJava 的flatMap把它压成一条直线:

api.getUser(userId) .flatMap(user -> api.getWeather(user.region)) .flatMap(weather -> api.getRecommendations(weather.season)) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe( products -> updateRecommendation(products), error -> handleError(error) );

flatMap的本质是“一对多映射”:它把上游User事件转换成一个新的Observable<Weather>,再把这个Observable的所有事件扁平化到当前流中。关键优势在于:错误传播是单向的、统一的。任何一个环节抛出异常(比如getWeather返回 404),整个链路会立即终止,跳转到最终的onError。你不需要在每一层写try-catch,也不用担心中间某个Callback忘记调用onFailure导致后续流程静默失败。我在重构一个老支付 SDK 时,把 17 个嵌套回调压缩成 3 个flatMap链,单元测试覆盖率从 32% 提升到 89%,因为每个环节的输入输出都变成了可验证的Observable类型。

3. 入门必须掌握的 5 个核心操作符:从“能跑”到“写得对”

3.1just()/fromCallable():不是玩具,是控制流的起点

很多教程一上来就用Observable.just(1, 2, 3),让人觉得 RxJava 就是“把数组包装一下”。其实just()的真实价值在于模拟同步操作的可观察性。比如你想测试一个处理用户登录状态的 ViewModel,但不想真调 API:

// 测试用:模拟登录成功 Observable<User> mockLogin = Observable.just(new User("test", "token123")); // 测试用:模拟登录失败 Observable<User> mockLoginFail = Observable.error(new IOException("Network error"));

fromCallable()才是生产环境的主力。它和just()的区别在于:just()是立即执行,fromCallable()是懒加载(lazy evaluation)。看这个例子:

// 错误:SharedPreferences 读取在创建 Observable 时就执行了! Observable<String> bad = Observable.just(sharedPrefs.getString("key", "")); // 正确:只有 subscribe 时才读取,避免在 Application 初始化时就触发 I/O Observable<String> good = Observable.fromCallable(() -> sharedPrefs.getString("key", ""));

fromCallable()返回的是ObservableSource,它的call()方法会在subscribe()时被调用。这保证了:1)I/O 操作不会在组件初始化时阻塞主线程;2)每次订阅都会获得最新值(比如用户在设置页修改了偏好,下次订阅就能拿到新值)。我在做离线优先 App 时,所有本地数据库查询都用fromCallable封装,配合cache()操作符,实现了“首次订阅查库,后续订阅直接返回缓存”的效果,比手写LruCache简洁 5 倍。

3.2map()/flatMap():数据转换的两种范式

map()是一对一转换,flatMap()是一对多展开。新手最容易混淆的是:什么时候该用flatMap?记住这个铁律:只要你的转换函数返回的是Observable(或SingleMaybe),就必须用flatMap。比如:

// 用户头像 URL → 下载图片 → Bitmap // 错误:map 返回 Observable,类型不匹配 Observable<Bitmap> wrong = Observable.just("https://...") .map(url -> downloadImage(url)); // downloadImage() 返回 Observable<Bitmap> // 正确:flatMap 展开 Observable Observable<Bitmap> correct = Observable.just("https://...") .flatMap(url -> downloadImage(url));

downloadImage(url)如果返回Observable<Bitmap>map()会把它当成普通对象塞进流里,导致下游收到的是Observable<Bitmap>而不是BitmapflatMap()则会订阅这个内部Observable,把它的所有onNext事件“压平”到外层流。另一个经典场景是“列表展开”:API 返回List<Item>,你想逐个处理每个Item

// 错误:map 把整个 List 当作一个元素 Observable<List<Item>> listStream = api.getItems(); listStream.map(items -> processItem(items.get(0))); // 只处理第一个 // 正确:用 flatMap + fromIterable 展开 listStream.flatMap(items -> Observable.fromIterable(items)) .map(item -> processItem(item)); // 每个 item 单独处理

这里fromIterable(items)List<Item>转成Observable<Item>flatMap再把它展开。实测下来,处理 1000 条数据时,flatMap+fromIterable比手写for循环快 15%,因为 RxJava 内部做了批量优化。

3.3filter()/takeUntil():事件流的“交通灯”

filter()很直观:保留满足条件的事件。但takeUntil()是真正的神器,它定义“流何时停止”。比如搜索框防抖:

// 每次输入触发 Observable<String> textChanges = RxTextView.textChanges(searchView); // 防抖:300ms 内无新输入才发出 Observable<String> debounced = textChanges .debounce(300, TimeUnit.MILLISECONDS); // 但 debounce 有个坑:用户快速输入 "abc" → "ab" → "a",debounce 会发出 "a"(最后一次) // 我们想要的是:只要用户开始输入,就取消上一次请求,只处理最后一次 Observable<String> latestOnly = textChanges .switchMap(query -> api.search(query) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) );

switchMapflatMap的升级版:它会取消前一个内部Observable的订阅,只保留最新的。takeUntil()则用于“主动截断”。比如监听 GPS 位置,直到用户点击“停止定位”按钮:

// 位置流 Observable<Location> locationStream = locationProvider.getLocationUpdates(); // 停止按钮点击流 Observable<Object> stopClicks = RxView.clicks(stopButton); // 只要 stopClicks 发出事件,locationStream 就停止 Observable<Location> activeLocations = locationStream.takeUntil(stopClicks);

takeUntil(stopClicks)的语义是:“locationStream发出的事件,只要stopClicks还没发出,就继续;一旦stopClicks发出,立刻 onComplete”。这比写if (isStopped) return;清晰 10 倍。

3.4distinctUntilChanged():避免重复渲染的隐形杀手

Android UI 最常见的性能陷阱:RecyclerView因为接收到重复数据而频繁notifyDataSetChanged()。比如一个开关按钮,用户猛点 5 次,isChecked()可能返回true, true, true, true, true。如果直接把isChecked()的结果映射到 UI:

// 危险:每次点击都触发更新,即使状态没变 RxCompoundButton.checkedChanges(switchButton) .map(isChecked -> getDisplayText(isChecked)) .subscribe(text -> textView.setText(text));

distinctUntilChanged()会对比相邻两个事件是否相等(用equals()),只在变化时转发:

RxCompoundButton.checkedChanges(switchButton) .distinctUntilChanged() // 连续 true 只发一次 .map(isChecked -> getDisplayText(isChecked)) .subscribe(text -> textView.setText(text));

它内部维护一个lastValue,每次onNext时比较current.equals(lastValue)。注意:distinctUntilChanged()默认用Object.equals(),如果你的User对象没重写equals(),它永远认为不相等。我在某社交 App 的个人资料页修复过这个问题:用户修改昵称后,头像、简介、标签三个模块同时刷新,但distinctUntilChanged()让它们只在真正变化时才重绘,帧率从 45fps 提升到 58fps。

3.5retryWhen():智能重试不是“循环 try-catch”

网络请求失败重试,新手常写:

// 错误:硬编码重试 3 次,无法区分错误类型 int retryCount = 0; while (retryCount < 3) { try { result = api.getData(); break; } catch (IOException e) { retryCount++; Thread.sleep(1000); } }

retryWhen()把重试逻辑声明式化:

api.getData() .retryWhen(errors -> errors .zipWith(Observable.range(1, 3), (error, retryCount) -> retryCount) .flatMap(retryCount -> { if (retryCount < 3) { return Observable.timer(1, TimeUnit.SECONDS); // 等待 1 秒 } else { return Observable.error(new RuntimeException("Max retry reached")); } }) );

errors是一个Observable<Throwable>,它发出所有上游的错误。zipWith把错误和重试次数配对,flatMap决定是否重试:返回Observable.timer()表示延迟重试,返回Observable.error()表示放弃。更高级的用法是按错误类型差异化重试

.retryWhen(errors -> errors .flatMap(error -> { if (error instanceof SocketTimeoutException) { // 超时:重试 2 次,间隔 2 秒 return Observable.range(1, 2) .flatMap(i -> Observable.timer(2, TimeUnit.SECONDS)); } else if (error instanceof HttpException && ((HttpException) error).code() == 401) { // 401:先刷新 token,再重试一次 return refreshToken().andThen(Observable.just(1)); } else { // 其他错误:不重试 return Observable.error(error); } }) );

这才是生产环境该有的重试策略:超时重试、认证失败走特殊流程、服务端错误直接上报。

4. 实操:从零搭建一个“搜索建议”功能——包含防抖、加载状态、错误处理全链路

4.1 项目结构准备:Gradle 依赖与基础配置

首先确认你的app/build.gradle已添加必要依赖。RxJava 3 是当前主流,它移除了Scheduler的静态方法(如Schedulers.io()),改用Schedules.io(),更符合函数式风格:

dependencies { // RxJava 核心 implementation 'io.reactivex.rxjava3:rxjava:3.1.6' // Android 线程调度器 implementation 'io.reactivex.rxjava3:rxandroid:3.1.0' // Retrofit 适配器(让 Call 变成 Observable) implementation 'com.squareup.retrofit2:adapter-rxjava3:2.9.0' // OkHttp 日志拦截器(调试用) debugImplementation 'com.squareup.okhttp3:logging-interceptor:4.12.0' }

注意版本对齐:rxjavarxandroidadapter-rxjava3必须同属 RxJava 3.x 系列。我踩过的坑是混用rxjava2rxandroid3,编译通过但运行时报NoSuchMethodError。另外,RxAndroidAndroidSchedulers.mainThread()依赖Looper.getMainLooper(),所以必须在主线程初始化。通常放在Application.onCreate()

public class MyApplication extends Application { @Override public void onCreate() { super.onCreate(); // 确保在主线程调用 AndroidSchedulers.init(); } }

4.2 API 接口定义:用 Retrofit + RxJava 3 封装

假设后端提供/search/suggestions?q={query}接口,返回 JSON:

{ "status": "success", "data": ["Android 开发", "RxJava 教程", "Jetpack Compose"] }

定义 Retrofit 接口:

public interface SearchApi { @GET("search/suggestions") Single<SearchResponse> getSuggestions(@Query("q") String query); } // 响应体 public class SearchResponse { public String status; public List<String> data; } // 创建 Retrofit 实例(单例) public class ApiClient { private static SearchApi searchApi; public static SearchApi getSearchApi() { if (searchApi == null) { OkHttpClient client = new OkHttpClient.Builder() .addInterceptor(new HttpLoggingInterceptor().setLevel(HttpLoggingInterceptor.Level.BODY)) .build(); Retrofit retrofit = new Retrofit.Builder() .baseUrl("https://api.example.com/") .client(client) .addConverterFactory(GsonConverterFactory.create()) .addCallAdapterFactory(RxJava3CallAdapterFactory.create()) // 关键:适配 RxJava 3 .build(); searchApi = retrofit.create(SearchApi.class); } return searchApi; } }

Single<SearchResponse>表示“最多发出一个事件(成功或失败)”,比Observable更语义化——搜索建议要么有结果,要么报错,不会发多个。

4.3 UI 层实现:搜索框 + RecyclerView + 加载状态

布局文件activity_search.xml

<LinearLayout xmlns:android="http://schemas.android.com/apk/res/android" android:layout_width="match_parent" android:layout_height="match_parent" android:orientation="vertical"> <com.google.android.material.textfield.TextInputLayout android:layout_width="match_parent" android:layout_height="wrap_content"> <com.google.android.material.textfield.TextInputEditText android:id="@+id/searchView" android:layout_width="match_parent" android:layout_height="wrap_content" android:hint="搜索..." /> </com.google.android.material.textfield.TextInputLayout> <ProgressBar android:id="@+id/progressBar" android:layout_width="wrap_content" android:layout_height="wrap_content" android:layout_gravity="center" android:visibility="gone" /> <androidx.recyclerview.widget.RecyclerView android:id="@+id/recyclerView" android:layout_width="match_parent" android:layout_height="0dp" android:layout_weight="1" /> </LinearLayout>

Activity 中的完整逻辑:

public class SearchActivity extends AppCompatActivity { private CompositeDisposable disposables = new CompositeDisposable(); private SearchAdapter adapter; @Override protected void onCreate(Bundle savedInstanceState) { super.onCreate(savedInstanceState); setContentView(R.layout.activity_search); EditText searchView = findViewById(R.id.searchView); ProgressBar progressBar = findViewById(R.id.progressBar); RecyclerView recyclerView = findViewById(R.id.recyclerView); adapter = new SearchAdapter(); recyclerView.setLayoutManager(new LinearLayoutManager(this)); recyclerView.setAdapter(adapter); // 核心:构建搜索流 Observable<String> queryStream = RxTextView.textChanges(searchView) .skip(1) // 跳过初始空字符串 .map(CharSequence::toString) .filter(text -> text.length() >= 2) // 至少 2 字才搜索 .debounce(400, TimeUnit.MILLISECONDS) // 防抖 400ms .distinctUntilChanged(); // 避免连续相同查询 // 处理搜索流 disposables.add( queryStream .switchMap(query -> { // 显示加载 progressBar.setVisibility(View.VISIBLE); // 发起网络请求 return ApiClient.getSearchApi() .getSuggestions(query) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .doOnError(error -> progressBar.setVisibility(View.GONE)) // 错误时隐藏进度条 .onErrorResumeNext(error -> { // 错误时返回空列表,避免崩溃 Toast.makeText(this, "搜索失败:" + error.getMessage(), Toast.LENGTH_SHORT).show(); return Single.just(new SearchResponse()); }); }) .subscribe( response -> { progressBar.setVisibility(View.GONE); adapter.submitList(response.data != null ? response.data : Collections.emptyList()); }, error -> { progressBar.setVisibility(View.GONE); Toast.makeText(this, "未知错误", Toast.LENGTH_SHORT).show(); } ) ); } @Override protected void onDestroy() { super.onDestroy(); // 生命周期绑定:取消所有订阅 disposables.clear(); } }

关键点解析:

  • skip(1)textChanges()EditText初始化时会发出一次空字符串,跳过它避免无效请求。
  • switchMap():用户快速输入时,自动取消上一次请求,只处理最后一次。
  • doOnError():在错误发生时隐藏进度条,这是副作用(side effect),doOnXXX系列操作符专门处理这类事。
  • onErrorResumeNext():把错误转换成正常事件(空列表),保证下游subscribe不会中断。这是 RxJava “错误恢复”的核心思想——错误不是终点,而是流的一部分。

4.4 Adapter 实现:用 ListAdapter + DiffUtil 优化 RecyclerView

SearchAdapter继承ListAdapter,利用DiffUtil自动计算差异:

public class SearchAdapter extends ListAdapter<String, SearchAdapter.ViewHolder> { public SearchAdapter() { super(DIFF_CALLBACK); } private static final DiffUtil.Callback DIFF_CALLBACK = new DiffUtil.Callback() { @Override public int getOldListSize() { return 0; // 由 ListAdapter 管理 } @Override public int getNewListSize() { return 0; } @Override public boolean areItemsTheSame(int oldItemPosition, int newItemPosition) { // 用索引判断,因为建议词可能重复 return oldItemPosition == newItemPosition; } @Override public boolean areContentsTheSame(int oldItemPosition, int newItemPosition) { // 内容相等才不刷新 List<String> oldList = getCurrentList(); List<String> newList = getCurrentList(); return Objects.equals(oldList.get(oldItemPosition), newList.get(newItemPosition)); } }; @NonNull @Override public ViewHolder onCreateViewHolder(@NonNull ViewGroup parent, int viewType) { View view = LayoutInflater.from(parent.getContext()) .inflate(android.R.layout.simple_list_item_1, parent, false); return new ViewHolder(view); } @Override public void onBindViewHolder(@NonNull ViewHolder holder, int position) { holder.textView.setText(getItem(position)); } static class ViewHolder extends RecyclerView.ViewHolder { TextView textView; ViewHolder(@NonNull View itemView) { super(itemView); textView = itemView.findViewById(android.R.id.text1); } } }

submitList()会触发DiffUtil计算,只刷新变化的 Item,比notifyDataSetChanged()性能高 3 倍。

4.5 单元测试:用 TestObserver 验证流逻辑

RxJava 的可测试性是它最大的优势之一。用TestObserver模拟整个链路:

@Test public void testSearchFlow_withValidQuery_emitsResults() { // 给定:mock API 返回固定数据 SearchApi mockApi = mock(SearchApi.class); when(mockApi.getSuggestions("rx")) .thenReturn(Single.just(new SearchResponse("success", Arrays.asList("RxJava", "RxAndroid")))); // 创建被测流 Observable<String> testStream = Observable.just("rx"); Observable<List<String>> resultStream = testStream .switchMap(query -> mockApi.getSuggestions(query) .map(response -> response.data) .onErrorReturn(throwable -> Collections.emptyList())); // 执行 TestObserver<List<String>> testObserver = resultStream.test(); // 验证 testObserver.assertValueCount(1); testObserver.assertValues(Arrays.asList("RxJava", "RxAndroid")); testObserver.assertComplete(); }

test()方法返回TestObserver,它提供了assertValueCount()assertValues()等断言方法,让你像测试普通方法一样测试异步流。

5. 常见问题排查与避坑指南:那些文档里不会写的实战经验

5.1 “内存泄漏”排查:CompositeDisposable不是万能的

CompositeDisposable.clear()确实能取消订阅,但它不能回收被 Observable 持有的 Activity 引用。看这个反模式:

// 危险:在 lambda 里捕获 this(Activity) Observable.interval(1, TimeUnit.SECONDS) .map(aLong -> { // 这里隐式持有 Activity.this updateUi(aLong); return aLong; }) .subscribe();

即使你调用了disposables.clear()interval产生的Observable依然持有Activity的强引用,直到interval自然结束(理论上永不结束)。正确做法是:所有 lambda 表达式里,只使用局部变量或静态方法

// 正确:用弱引用或提取方法 private void setupTimer() { disposables.add( Observable.interval(1, TimeUnit.SECONDS) .map(this::updateUiFromTimer) // 提取为实例方法 .subscribe() ); } private Long updateUiFromTimer(Long aLong) { // 在这里访问 UI 组件 if (!isFinishing() && !isDestroyed()) { textView.setText("Time: " + aLong); } return aLong; }

或者用WeakReference显式管理:

WeakReference<SearchActivity> activityRef = new WeakReference<>(this); disposables.add( Observable.interval(1, TimeUnit.SECONDS) .map(aLong -> { SearchActivity activity = activityRef.get(); if (activity != null && !activity.isFinishing()) { activity.updateUi(aLong); } return aLong; }) .subscribe() );

5.2 “线程切换失效”:subscribeOn()observeOn()的作用域陷阱

新手常以为subscribeOn()会“全局生效”,其实它只影响上游操作符的执行线程。看这个错误示例:

// 错误:observeOn 在 map 之后,但 map 里的耗时操作仍在 io 线程执行! api.getData() .subscribeOn(Schedulers.io()) .map(data -> { // 这里是耗时解析(如 JSON 解析),但它在 io 线程执行 return parseJson(data); }) .observeOn(AndroidSchedulers.mainThread()) .subscribe(result -> updateUi(result));

parseJson()io线程执行,虽然updateUi()在主线程,但解析过程依然阻塞io线程池。正确做法是:把耗时操作也放到subscribeOn()指定的线程

// 正确:用 flatMap 把解析也移到 io 线程 api.getData() .subscribeOn(Schedulers.io()) .flatMap(data -> Observable.fromCallable(() -> parseJson(data))) .observeOn(AndroidSchedulers.mainThread()) .subscribe(result -> updateUi(result));

或者,如果parseJson()很快,就别管它;如果很慢,用fromCallable包装。

5.3 “空指针”高频场景:onError里调用getView()的陷阱

onError回调一定在observeOn指定的线程执行,但它不保证 Activity 还活着。用户可能在请求过程中按了返回键:

// 危险:onError 里直接调用 findViewById .subscribe( result -> updateUi(result), error -> { // 此时 Activity 可能已 destroy,findViewById 返回 null Toast.makeText(this, error.getMessage(), Toast.LENGTH_SHORT).show(); } );

解决方案是:onError里检查 Activity 状态

error -> { if (!isFinishing() && !isDestroyed()) { Toast.makeText(this, error.getMessage(), Toast.LENGTH_SHORT).show(); } }

更健壮的做法是封装一个SafeToast工具类:

public class SafeToast { public static void show(Context context, String message) { if (context instanceof Activity) { Activity activity = (Activity) context; if (!activity.isFinishing() && !activity.isDestroyed()) { Toast.makeText(activity, message, Toast.LENGTH_SHORT).show(); } } else { Toast.makeText(context, message, Toast.LENGTH_SHORT).show(); } } }

5.4 “背压溢出”:Observable.create()的致命错误

Observable.create()是最危险的操作符,它不处理背压(backpressure)。比如:

// 危险:无限发射,没有背压控制 Observable.create(emitter -> { while (true) { emitter.onNext(System.currentTimeMillis()); // 每毫秒发一个 } });

下游如果消费不过来(比如observeOn(AndroidSchedulers.mainThread())),就会 OOM。正确做法是:Flowable替代Observable,并指定背压策略

Flowable<Long> flowable = Flowable.create(emitter -> { while (!emitter.isCancelled()) { emitter.onNext(System.currentTimeMillis()); Thread.sleep(100); // 降低发射频率 } }, BackpressureStrategy.LATEST); // LATEST:只保留最新一个

BackpressureStrategy选项:

  • MISSING:不处理背压(同Observable,危险)
  • ERROR:下游来不及消费时抛异常
  • BUFFER:缓存所有事件(可能 OOM)
  • LATEST:只保留最新一个(推荐)
  • DROP:新事件到来时丢弃旧事件

5.5 “调试困难”:如何用doOnSubscribe()doOnEvent()定位问题

RxJava 链路长,出问题很难定位。doOnXXX系列是调试利器:

api.search(query) .doOnSubscribe(disposable -> Log.d("SEARCH", "开始搜索:" + query)) .doOnNext(result -> Log.d("SEARCH", "收到结果:" + result.size())) .doOnError(error -> Log.e("SEARCH", "搜索失败", error)) .doOnComplete(() -> Log.d("SEARCH", "搜索完成")) .subscribe(...);

doOnSubscribe()subscribe()被调用时触发,doOnNext()在每次onNext()前触发。它们不改变流,只打日志。我在线上环境用doOnNext()发现过一个 bug:后端返回的List里有null元素,导致RecyclerView绑定时 NPE,而doOnNext()日志清晰显示了null的存在位置。

6. 进阶思考:RxJava 在现代 Android 架构中的位置——它会被 Jetpack Compose 取代吗?

6.1 Compose 的StateFlowSharedFlow:相似但不同

Compose 推荐用StateFlow(热流)和SharedFlow(冷流)替代 RxJava。它们和 RxJava 的对应关系:

RxJavaCompose Flow特点
ObservableSharedFlow多对多,支持重放,但不保证顺序
SingleStateFlow一对一,始终有值,UI 自动重组
SubjectMutableStateFlow可变的 StateFlow

关键区别在于:Flow 是 Kotlin 协程原生的,而 RxJava 是 Java 的。这意味着:

  • Flow 可以用collectLatest实现switchMap的效果:
    searchQuery.collectLatest { query -> api.search(query).collect { results -> searchState.value = results } }
  • Flow 的错误处理更 Kotlin 化:catch块直接throw,无需onErrorResumeNext
  • Flow 的生命周期
http://www.cnnetsun.cn/news/4187203.html

相关文章:

  • 完整跑通 tmom 多厂区 MOM/MES 系统:从部署到车间过站的实操手册
  • 从流程图到状态机:嵌入式开发中的事件驱动编程范式
  • Java面试实战:技术深度与软素质双维度考察
  • Java后端面试核心:SQL优化、HashMap并发与内存调优
  • 从OpenClaw到Hermes:AI智能体开发工具链的升级与实战迁移指南
  • 2026年Java面试题库:核心考点与趋势解析
  • 多智能体与领域知识驱动的代码适配框架:从Spring Boot到Quarkus的自动化迁移实践
  • 链表数据结构与面试核心要点解析
  • Python win32com自动化Office与Outlook:从原理到实战报表邮件系统
  • 电力约束下数据中心转型:从算力军备竞赛到能效优化实战
  • 算法日常・每日刷题--<BFS最短路径>4
  • 深入解析RS232、RS422、RS485串口通信:从电气原理到工业应用实战
  • Hermes Agent 日志监控系统搭建教程:ELK 一键部署 + 智能异常检测完整指南
  • 碧蓝航线自动化指南:5分钟配好 Alas,日常全托管
  • 文件包含漏洞实战:从CTF赛题看PHP特性与LFI2RCE利用链
  • 27考研408操作系统强化课程:高效攻克进程管理与内存管理核心考点
  • 开源框架WithEveryone:解决多角色图像生成的身份一致性与场景规划难题
  • 机器人百米冲刺与替代人工:核心技术解析与ROS仿真实践
  • 2026年软件测试面试高频考点与实战策略
  • Windows驱动开发:自签名证书原理与实战,解决驱动强制签名问题
  • FOC控制核心数学工具:正余弦查找表、Atan2与限幅的嵌入式实现
  • 树莓派无头启动SSH连接全攻略:四种方法获取IP与深度排错
  • MATLAB浮点转定点实战:Q格式量化与硬件部署避坑指南
  • CursorRules 实战指南:3 步让 AI 助手写出符合你项目规范的代码
  • Flash浏览器CefFlashBrowser:5分钟救活你的SWF老游戏
  • SpringBoot实习管理系统架构设计与实践
  • 《OPC智能体:一个人的容度智能体》白皮书——专知智库OPC研究院关于“岗位级智能体”的官方定义与产业实践白皮书
  • FreeRTOS任务通知在STM32上的底层原理与实战应用
  • 基于MinerU为Claude Code构建本地PDF解析技能,实现文档智能处理
  • 从GitHub中断看被动扩展瓶颈:高可用架构的主动防御策略