Как правильно реализовать цепочку rxjava в android kotlin?
Я пытаюсь реализовать условия отсутствия интернета на мобильном устройстве. При этом должна происходить такая последовательность операций в цепочке:
1.Проверяем время последнего запроса
2.Делаем запрос на rest api
3.В случае удачи мы наполняем кэш и БД данными с сервера, в противном случае(он работает неверно) мы отлавливаем ошибку, берем данные из БД и наполняем ими PublisherSubject. Потом мы включаем связь и дергаем Swipe fefresh нашего фрагмента. Я пробовал методы doOnError onErrorReturn onErrorResumeNext. В каждом из них происходила такая интересная ситуация в коде: сначала дергается (1) потом (2) и потом, не могу понять почему (3), из-за чего данные приходят несколько раз, с каждым проходом увеличиваются на 1, как будто цепочка накапливает ошибки. Кто может с этим сталкивался и может помочь мне с этой проблемой? Буду очень благодарен
Так же в цепочке присутствует внешняя зависимость от Системного времени, что не хорошо, если по ней будут какие-то статьи на прочтение тоже буду признателен
override fun requestAllCountries(): Flowable<Any> {
return getTimeSinceLastUpdate()
.flatMap { isMoreThanMinute ->
return@flatMap if (isMoreThanMinute) {
lastNetworkRequestTime = System.currentTimeMillis()
networkRepository.getCountryDate()
.doOnError {
(3) databaseRepository.getAllCountries().subscribe({countrySubject.onNext(it)}, {})
}
.flatMap {
it.forEach { item ->
databaseRepository.addLanguage(item.languages)
}
(1) databaseRepository.addAllCountries(it)
(2) cacheRepository.addAllCountries(it)
}
} else {
cacheRepository.getAllCountries()
}
}
.doOnNext {
Log.e("HZ", "$it")
countrySubject.onNext(it)
}
.map { Any() }
}
Ответы (1 шт):
Судя по всему, у вас основная проблема - в подписывании на источник данных внутри другого источника данных. Что, конечно, можно решить управляя завершением подписки. Однако, лучше вообще избегать таких решений - подписки внутри источника.
Алгоритм проблемный у вас вот почему:
- Выключили сеть.
- Запустили запрос в сеть.
- Подписались на изменения из БД. Оно у вас, видимо, не Single, а Flowable(или Observable). Т.е. выдаёт событие при изменениях в БД, дергая Subject.
- После включения сети вы получаете данные с сервера.
- Пришедшие данные записываете в БД, вызывая срабатывание ранее зарегистрированных подписок на эти изменения (см. п.3)
Простым (и очень сомнительным) решением будет заменить тип источника данных databaseRepository.getAllCountries() на Single. Так он сработает всего один раз и не будет реагировать на дальнейшие изменения в БД.
Правильным решением же было бы не создание внутренней подписки, но модификация имеющегося источника данных. Главное - уследить какие типы источников используются. Мне неизвестен ожидаемый принцип работы, так что предположу, что вам надо просто получить данные из кэша или сети (записав данные в БД и кэш). И не надо, чтобы источник данных продолжал жить своей сложной жизнью:
override fun requestAllCountries(): Completable {
return getTimeSinceLastUpdate() // должен быть Single
.flatMap { isMoreThanMinute ->
if (isMoreThanMinute) {
lastNetworkRequestTime = System.currentTimeMillis()
networkRepository.getCountryDate() // должен быть Single
.doOnSuccess {
it.forEach { item ->
databaseRepository.addLanguage(item.languages) // не должно быть Rx
}
databaseRepository.addAllCountries(it) // не должно быть Rx
cacheRepository.addAllCountries(it) // не должно быть Rx
}
.flatMap {
databaseRepository.getAllCountries() // должен быть Single
}
.onErrorResumeNext(
databaseRepository.getAllCountries() // должен быть Single
)
} else {
cacheRepository.getAllCountries() // должен быть Single
}
}
.doOnSuccess {
Log.e("HZ", "$it")
countrySubject.onNext(it)
}
.ignoreElement()
}