RxJava列表转换异常:Flowable<List<TaskEntity>>转Flowable<List<Task>>无输出
我最近在把Flowable<List<TaskEntity>>转换成Flowable<List<Task>>的时候踩了个坑——测试用的简单列表转换完全正常,但一用到DAO返回的数据流就没输出,折腾了好一会儿才找到原因,分享给大家:
测试场景:逻辑验证正常
我先写了个简单的测试代码验证转换逻辑,这段是能正常输出预期结果的:
Flowable.fromArray(Arrays.asList(1,2,3)) .flatMapIterable(ids->ids) .map(s->"No. "+s) .toList() .toFlowable() .subscribe( t -> Log.d(TAG, "getAllActiveTasks: "+t) );
输出结果符合预期:[No.1 No.2 No.3]
实际场景:无输出的问题代码
但把同样的逻辑用到DAO返回的数据流时,却完全没有输出:
mTaskDao.getAllTasks(STATE_ACTIVE) .flatMapIterable(task -> task) .map(Task::create) .toList() .toFlowable() .subscribe( t -> Log.d(TAG, "getAllActiveTasks: "+t) );
附Task.create的实现:
public static Task create(TaskEntity eTask) { Task task = new Task(eTask.getTaskId(), eTask.getTaskTitle(), eTask.getTaskStatus()); task.mTaskDescription = eTask.getTaskDescription(); task.mCreatedAt = eTask.getCreatedAt(); task.mTaskDeadline = eTask.getTaskDeadline(); return task; }
问题根源:toList()对无限数据流无效
后来才搞明白,问题出在toList()这个操作符上!toList()的核心特性是必须等上游数据流完全结束(发送onComplete事件),才会把收集到的所有元素打包成列表发送出去。
测试用的Flowable.fromArray(...)是有限数据流——它发送完列表里的元素后就会立刻发送onComplete,所以toList()能正常收集并输出结果。但DAO(比如Room的DAO)返回的Flowable是无限数据流:它会持续监听数据库的变化,只要数据有更新就会重新发送新的列表,永远不会主动发送onComplete事件。这就导致toList()一直等着上游结束,永远收不到完整的列表,自然就没有输出了。
解决方案:直接转换列表或处理单个元素后局部收集
既然知道了原因,解决起来就简单了,这里提供两种常用方案:
方案1:直接转换整个列表(推荐)
我们不需要把列表拆成单个元素再收集回来,直接对上游的每个列表做转换就好,代码更简洁高效:
mTaskDao.getAllTasks(STATE_ACTIVE) .map(taskEntities -> taskEntities.stream() .map(Task::create) .collect(Collectors.toList()) ) .subscribe(t -> Log.d(TAG, "getAllActiveTasks: "+t));
这样一来,每次上游发送新的List<TaskEntity>,我们就直接把它转换成List<Task>并发送,既保留了数据库监听的特性,又完成了类型转换。
方案2:处理单个元素后局部收集
如果你的业务逻辑必须先处理单个TaskEntity(比如要做过滤、复杂的单元素转换),可以把每个上游发来的列表单独处理成一个有限数据流,再用toList()收集:
mTaskDao.getAllTasks(STATE_ACTIVE) .flatMap(taskEntities -> Flowable.fromIterable(taskEntities) .map(Task::create) // 这里可以添加单个元素的处理逻辑,比如filter、doOnNext等 .toList() .toFlowable() ) .subscribe(t -> Log.d(TAG, "getAllActiveTasks: "+t));
这个写法里,内部的Flowable.fromIterable(taskEntities)是有限数据流(处理完列表元素就会发送onComplete),所以toList()能正常工作,处理完后再把转换好的列表发送出去。
内容的提问来源于stack exchange,提问作者Sahil Patel

