Страницы

Поиск по вопросам

Показаны сообщения с ярлыком rxjava. Показать все сообщения
Показаны сообщения с ярлыком rxjava. Показать все сообщения

вторник, 31 марта 2020 г.

Как грамотно отправить 2-3 запроса android

#java #android #rxjava


Стоит задача, отправить пару запросов, получить с этих запросов данные, обработать
их (данные из 2-3 запросов) и отправить в адаптер( например)

Каким образом это правильно и красиво сделать?
Насколько я знаю, для таких задач используется RxJava, хотелось бы чтобы привели
пример правильного решение с использованием Rx и без использование Rx.
    


Ответы

Ответ 1



Если вам нужно соеденить ответы этих запросов то вполне подойдет оператор concat Вот небольшая демонстрация final String[] aStrings = {"A1", "A2", "A3", "A4"}; final String[] bStrings = {"B1", "B2", "B3"}; final Observable aObservable = Observable.fromArray(aStrings); final Observable bObservable = Observable.fromArray(bStrings); Observable.concat(aObservable, bObservable) .subscribe(getObserver()); На выходе получим поток данных спрерва "A1", "A2", "A3", "A4" а затем "B1", "B2", "B3" Так же можно использовать оператор merge Observable.merge(aObservable, bObservable) .subscribe(getObserver()); тогда на выходе получим поток данных "A1", "B1", "A2", "A3", "A4", "B2", "B3" или же zip Observable stringObservable1 = Observable.just("Hello", "World"); Observable stringObservable2 = Observable.just("Bye", "Friends"); Observable.zip(stringObservable1, stringObservable2, new BiFunction() { @Override public String apply(@NonNull String s, @NonNull String s2) throws Exception { return s + " - " + s2; } }).subscribe(new Consumer() { @Override public void accept(String s) throws Exception { System.out.println(s); } }); на выход получим Hello - Bye World - Friends Я б не стал писать это чисто андроид потоками ибо это ресурсаемко советую использовать rxJava.

пятница, 13 марта 2020 г.

AlertDialog с помощью RxJava

#android #rxjava


С помощью AsyncTask при обработке массива данных мы легко можем вывести с помощью
AlertDialog прогресс выполнения обработки.

Можно ли это сделать с помощью RxJava? Прошу вашего простейшего примера.

И как решается проблема поворота девайса при загрузке данных с сервера при помощи
RxJava и отображения прогресса на AlertDialog?
    


Ответы

Ответ 1



Самый простой способ это заюзать метод from: List list = new ArrayList<>(); for (int i = 0; i < 100; i++) { list.add("item "+i); } Observable.from(list) .observeOn(AndroidSchedulers.mainThread()) .subscribe(new Subscriber() { @Override public void onCompleted() { if(progress!=null && progress.isShowing()) progress.dismiss(); } @Override public void onError(Throwable e) { if(progress!=null && progress.isShowing()) progress.dismiss(); } @Override public void onNext(String s) { progress = ProgressDialog.show(this, "dialog title", "dialog message", true); } });

воскресенье, 8 марта 2020 г.

java.lang.IllegalStateException: Expected BEGIN_ARRAY but was BEGIN_OBJECT (Retrofit 2)

#java #android #retrofit #rxjava


Получаю такую ошибку java.lang.IllegalStateException: Expected BEGIN_ARRAY but was
BEGIN_OBJECT, когда получаю данные с сервера. Чем вызвана эта ошибка? И как ее исправить?  

JSON, который приходит с сервера: 

{
  "group": [
    {
      "name": "Group1",
      "description": "Test Group 1"
    },
    {
      "name": "Group2",
      "description": "Group Name Updated"
    }
  ]
}


Интерфейс API: 

 @GET("groups")
    Observable > getAllGroups(@Header("Authorization") String auth,
                                   @Header("Content-type") String contentType,
                                   @Header("Accept") String accept
                                  );


Метод, в котором получаю данные: 

private void getAllGroups() {
    String group = "Group1";
    String credentials = "admin" + ":" + "admin";
    final String basic =
            "Basic " + Base64.encodeToString(credentials.getBytes(), Base64.NO_WRAP);
    String contentType = "application/json";
    String accept = "application/json";
    Subscription subscription = App.service.getAllGroups(basic, contentType, accept)
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .subscribe(groups -> {
                groupList.addAll(groups);
            }, throwable -> {
                Log.e("All group error", String.valueOf(throwable));
            });

    addSubscription(subscription);
}


Класс Group: 

public class Group {

    private String name;
    private String description;
    private String admins;
    private Member members;

    public Group(String name, String description, String admins, Member members) {
        this.name = name;
        this.description = description;
        this.admins = admins;
        this.members = members;
    }

    public Group(String name, String description) {
        this.name = name;
        this.description = description;
    }
  // getters setters
}

    


Ответы

Ответ 1



Сервер присылает JSONObject, в котором JSONArray, а вы пытаетесь парсить как JSONArray. Должно быть как то так: public class GroupList { ArrayList group; //setters, getters } public class Group { private String name; private String description; ... }

среда, 4 марта 2020 г.

Методы just в RxJava

#java #android #rxjava


Просматривая исходники библиотеки io.reactivex.rxjava2, в  классе Observable обнаружила
интересную (или странную?) вещь - 10 методов just, сигнатура которых выглядит следующим
образом:

just(T item)
just(T item1, T item2)
just(T item1, T item2, T item3)
just(T item1, T item2, T item3, T item4)
just(T item1, T item2, T item3, T item4, T item5)
just(T item1, T item2, T item3, T item4, T item5, T item6) 
just(T item1, T item2, T item3, T item4, T item5, T item6, T item7)
just(T item1, T item2, T item3, T item4, T item5, T item6, T item7, T item8)
just(T item1, T item2, T item3, T item4, T item5, T item6, T item7, T item8, T item9) 
just(T item1, T item2, T item3, T item4, T item5, T item6, T item7, T item8, T item9,
T item10) 


В чём глубокий смысл этого?.. Почему нельзя было сделать как-нибудь вроде just(T...
item)? Почему методов именно 10?
    


Ответы

Ответ 1



Вероятно, это сделано, чтобы эти статические методы можно было приводить к функциональным интерфейсам семейства rx.functions.Func* public interface Func1 extends Function { R call(T t); } ... public interface Func9 extends Function { R call(T1 t1, T2 t2, T3 t3, T4 t4, T5 t5, T6 t6, T7 t7, T8 t8, T9 t9); }

понедельник, 2 марта 2020 г.

Как превратить EditText в поток данных для RxJava2?

#android #rxjava #listener #rxandroid #rxjava2


Начал постепенно переходить на реактивное программирование с замечательным фреймворком
RxJava 2 и заинтересовал простейший пример создания потока в динамическом стиле. К
сожалению, не считаю пример Observable.just(1,2,3) достаточно полным для себя, но вот
уже несколько часов мучаю себя мыслью, как превратить изменения в EditText в поток
данных? Если не сложно, опишите, пожалуйста, пример как сделать реальный Observable
из TextWatcher, чтобы можно было на него подписаться и подписчик реагировал на изменения
текста в EditText. 
    


Ответы

Ответ 1



Можно воспользоваться готовой библиотекой RxBinding. Также в ней можно посмотреть конкретную реализацию. Вообще для связки callback-методов с rx используется метод Observable.create, например: Observable.create(emitter -> { TextWatcher watcher = new TextWatcher() { @Override public void beforeTextChanged(CharSequence charSequence, int i, int i1, int i2) { } @Override public void onTextChanged(CharSequence charSequence, int i, int i1, int i2) { } @Override public void afterTextChanged(Editable editable) { if (!emitter.isDisposed()) { //если еще не отписались emitter.onNext(editable.toString()); //отправляем текущее состояние } } }; emitter.setCancellable(() -> editText.removeTextChangedListener(watcher)); //удаляем листенер при отписке от observable editText.addTextChangedListener(watcher); }); Нужно не забыть подписаться, и самое главное, отписаться от такого источника данных, т.к. он держит ссылку на editText и может привести к утечке памяти.

воскресенье, 12 января 2020 г.

RxJava не работает с файлами в io() потоке

#java #rxjava


Суть проблемы такова: у меня есть Observable и Subscriber. Observable я пытаюсь запустить
в io() потоке, так как он работает с файлами (не буду показывать код, он довольно большой,
но не имеет значения), однако он ничего не делает:

    Observable creatingObservable = getCreatingObservable(image);
    Subscriber creatingSubscriber = getCreatingSubscriber();

    creatingObservable
            .subscribeOn(Schedulers.io())
            .subscribe(creatingSubscriber);


Если запускать код без subscribeOn - все прекрасно работает. Так в чем же проблема
и как ее исправить?



P.S. У меня еще System.out.println() не работает. Проблема распространяется на все
потоки Scheduler'a.
    


Ответы

Ответ 1



Проблема была в том, что главный поток не дожидался окончания окончания выполнения RxJava'вского потока. В результате, RxJava даже не успевал "пискнуть" - от сюда никаких сообщений из System.out.println() и работы с файлами. Решение подсказали тут - https://stackoverflow.com/questions/37993371/rxjava-doesnt-work-in-scheduler-io-thread

суббота, 11 января 2020 г.

Как выполнить более 3 параллельных http запросов с помощью RxJava (+retrofit2) в Android

#android #kotlin #rxjava #retrofit2


Пытаюсь освоить работу с RxJava + retrofit2 + kotlin.

Как в этом участвует Observable в общем разобрался, даже научился объединять результаты
2-3 параллельных запросов по итогу выполнения всех с помощью Observable.zip, НО как
похожим образом выполнить более 3 параллельных запросов не могу понять.

В исходном коде класса Observable есть метод 

static  Observable zip(), 


но он не работает как 

static  Observable zip() 


в случае с тремя Observable на вход и Function3<*,*,*,*>. 

Старался уже сделать первые три запроса с помощью zip, потом flatMap и тд, но всё
не выходит. Чтение документации и примеров не помогают повернуть к нужному направлению.
Смотрел в сторону Observable.combineLatest, но пришел к выводу (может ошибочному),
что метод вернет результат первого выполнившегося Observable.

Если в целом, что я хочу: 

У меня есть 4 метода, которые возвращают Observable:


this.orderRepository.getStatuses():Observable
this.orderRepository.getOrders():Observable
this.userRepository.getUsers():Observable
this.orderRepository.getTypes():Observable


Подскажите, пожалуйста, как я должен скомпоновать эти методы, чтобы они выполнились
параллельно, и в конце получить результат выполнения всех 4 методов в одном месте? 

Мои стыдливые попытки ниже. Это уже что-то испорченное, но я вообще не понимаю, что
я должен сделать с 4 Observables, чтобы их всех в одном месте и вернуть в UI. + уже
запутался с flatMap, map и так далее...

override fun buildUseCaseObservable(params: Params): Observable {

return Observable.zip(
                this.userRepository.getUsers(100, 1, ApiUserFilter(isManager = true)),
                this.orderRepository.getOrderStatuses(),
                this.orderRepository.getOrderTypes(),
                Function3
{ t1, t2, t3 ->
                    Zip(t1,t2,t3)
                }).map {
                    this.orderRepository.getOrders(params.limit, params.page, params.filter).flatMap {
                        it
                    }
                }

 }



    


Ответы

Ответ 1



Zip для четырех observable можно написать так: Observable.zip( getStatuses(), getOrders(), getUsers(), getTypes(), Function4 { status, order, user, type -> ... } )

четверг, 9 января 2020 г.

Как правильно заменить AsyncTask с помощью RxJava2 и RxAndroid2

#android #rxjava #rxandroid #rxjava2


У меня есть AsyncTask, который я хочу переделать в Rx 

Вот так выглядит мой AsyncTask

new AsyncTask()
    {
        ArrayList undoneServiceCodes = new ArrayList<>();
        HashMap undoneForms = new HashMap<>();
        boolean isHasAtLeastOneDoneServiceCode = false;

        @Override
        protected void onPreExecute()
        {
            super.onPreExecute();

            if (iCloseCallListener != null)
            {
                iCloseCallListener.onPreValidation();
            }
        }

        @Override
        protected Void doInBackground(Void... params)
        {
            undoneServiceCodes = getUndoneServiceCodes();
            undoneForms = getUndoneForms();
            isHasAtLeastOneDoneServiceCode = isHasAtLeastOneDoneServiceCode();

            return null;
        }

        @Override
        protected void onPostExecute(Void iVoid)
        {
            super.onPostExecute(iVoid);

            if (iCloseCallListener != null)
            {
                iCloseCallListener.onPostValidation(undoneServiceCodes, undoneForms,
isHasAtLeastOneDoneServiceCode);
            }
        }
    }.execute();


Мне не понятно как можно выполнить в бекграунде Rx , можно выполнить сразу 3 разных
метода как в примере с AsyncTask когда в бекграунде выполняется 3 метода


undoneServiceCodes = getUndoneServiceCodes();
undoneForms = getUndoneForms();
isHasAtLeastOneDoneServiceCode = isHasAtLeastOneDoneServiceCode();


Если это был бы один метод(допустим первый) я бы это сделал так

Flowable.fromIterable(getUndoneServiceCodes())//
            .subscribeOn(Schedulers.io())//
            .observeOn(AndroidSchedulers.mainThread())//
            .doOnSubscribe(iSubscription ->
            {
                if (iCloseCallListener != null)
                {
                    iCloseCallListener.onPreValidation();
                }
            }).toList()//
            .subscribe(resultList -> {
                if (iCloseCallListener != null)
                {
                    iCloseCallListener.onPostValidation(iCloseCallListener, ???, ???);
                }
            });


Но так я получу результат только для одного выполняемого в бекграунде метода, как
сделать так, чтоб можно было обрабоать 3?
    


Ответы

Ответ 1



Можно воспользоваться оператором zip как-то так (точность названий методов не гарантирую): Flowable.zip( Flowable.fromCallable(method1()), Flowable.fromCallable(method2()), Flowable.fromCallable(method3()), (result1, result2, result3) -> new Triple(result1, result2, result3) ) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(triple -> iCloseCallListener.onPostValidation(triple.first, triple.second, triple.third))

пятница, 3 января 2020 г.

Обработка множественных запросов на RxJava

#java #android #многопоточность #rxjava


Есть поле для вводе текста, допустим ищем пользователей по имени в БД.
Есть метод (упрощенный), который выполняется при каждом наборе символа.

    public void loadUsers(){
        getUsers(mName)
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .subscribe(users -> {
             //отображение данных
            }, throwable -> {
             //обработка ошибки
            });
}


Проблема данного метода в том, что при быстром наборе текста вполне вероятна ситуация,
когда ответ на предпоследний запрос пришел позже, чем на последний. Тогда отображаемые
данные будут не соответствовать введенному тексту.
Вопрос: как обработать запросы в том порядке, в котором выполнялись?
Или, может быть, можно игнорировать предыдущие запросы? 
P.S. С использованием RxJava.
    


Ответы

Ответ 1



Как вариант можно хранить Subscription и её сначала обнулять, а потом новый запрос слать. Типа как-то так (за названия методов не ручаюсь): Subscription s; public void loadUsers() { if(s != null && !s.isUnsubscribed()){ s.unsubscribe(); } s = getUsers(mName) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(users -> { //отображение данных }, throwable -> { //обработка ошибки }); }

Ответ 2



Для того чтобы легко обработать эту ситуацию можете использовать библиотеку RxBinding. Она переводит слушатели на Rx что даёт нам легко воспользоваться оператором debounce(). Почитать про оператор можно здесь. Должно получиться примерно так: public void loadUsers(){ RxTextView.afterTextChangeEvents(mNameEditText) .debounce(300, TimeUnit.MILLISECONDS) .map(TextViewAfterTextChangeEvent::editable) .map(Editable::toString) .flatMap(name -> getUsers(name)) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(users -> { //отображение данных }, throwable -> { //обработка ошибки }); }

Правильное выполнение запроса RxJava+Retrofit

#android #rest #retrofit #rxjava


Здравствуйте. Недавно начал разбираться с RxJava. Что-то тяжело пока дается...
Есть сервер на котором есть приватные чаты(тет-а-тет) и групповые.
Посредством Rest запроса на сервер нужно выудить следующую инфу оттуда:
Все существующие чаты и диалоги:

Observable> chats = apiService.getChats();
Observable> dialogs = apiService.getDialogs();


Объекты Chat и Dialog содержат переменные:

int unreadMessagesCount (количество непрочитанных сообщений);
int id (по этому id запрашивается список сообщений из чата)


Запросы на список сообщений

Observable> dialogMessages = apiService.getDialogMessages(String id);
Observable> chatMessages = apiService.getChatMessages(String id);


Как в этом случае более грамотно составить запросы используя RxJava, чтобы опросить
все чаты и диалоги на предмет новых сообщений и затем получить эти сообщения в один список?
    


Ответы

Ответ 1



Например как-то так: Получаем массив чатов. Преобразуем массив оных в очередь объектов Chat. Получаем детали каждого. Результат преобразовываем обратно в массив:. apiService.getChats() .from(Observable::from) .flatMap(chat -> apiService.getChatMessages(chat.id)) .toList() .subscribe(System::out); Если без лямбд, то from(Observable::from) можно переписать вот такой ужасной конструкцией: .flatMap(new Func1, Observable>() { @Override public Observable call(List chats) { return Observable.from(chats); } })

пятница, 20 декабря 2019 г.

Различие между subscribeOn и observeOn методами

#rxjava


В чём заключается различие между subscribeOn и observeOn методами?

Правильно ли я его понимаю? 

Вот то, как я это понимаю: subscribeOn определяет поток по умолчанию для Observable
после его создания (в случае, если его нужно выполнять не в текущем потоке), т.о.,
начинаться выполнение будет всегда в потоке, определённом subscribeOn. И поэтому subscribeOn
нужен только один (если будет несколько subscribeOn, выполнится только первый). А observeOn
после может поменять поток, начиная с места вызова, и сколько будет observeOn, столько
раз будет меняться поток.
    


Ответы

Ответ 1



Нашла полезную статью, которая содержит точный ответ насчёт количества и места вызова данных методов. Надеюсь, поможет кому-нибудь ещё. По поводу subscribeOn() The subscribeOn() operator will have the same effect no matter where you place it in the observable chain; however, you can't use multiple subscribeOn() operators in the same chain. If you do include more than one subscribeOn(), then your chain will only use the subscribeOn() that’s the closest to the source observable. Что значит: Оператор subscribeOn() будет иметь тот же эффект независимо от того, где вы поместите его в цепочку observable; однако вы не можете использовать несколько операторов subscribeOn() в одной цепочке. Если вы включили в цепочку более одного subscribeOn(), ваша цепочка будет использовать только subscribeOn(), который ближе всего к observable источнику. subscribeOn(), который ближе всего к Observable источнику - это первый в цепочке, и данное утверждение было также проверено опытным путём, например, если есть цепочка: observable .subscribeOn(Schedulers.newThread()) .subscribeOn(Schedulers.computation()) .subscribe(observer); то onNext() каждого элемента выполняется в потоке RxNewThreadScheduler-1 (потоке, созданном Schedulers.newThread()), а цепочка observable .subscribeOn(Schedulers.computation()) .subscribeOn(Schedulers.newThread()) .subscribe(observer); передаёт выполнение onNext() каждого элемента в поток RxComputationThreadPool-1 (поток, созданном Schedulers.computation()). По поводу observeOn() Unlike subscribeOn(), where you place observeOn() in your chain does matter, as this operator only changes the thread that’s used by the observables that appear downstream. For example, if you inserted the following into your chain then every observable that appears in the chain from this point onwards will use the new thread. .observeOn(Schedulers.newThread()) This chain will continue to run on the new thread until it encounters another observeOn() operator, at which point it’ll switch to the thread specified by that operator. You can control the thread where specific observables send their notifications by inserting multiple observeOn() operators into your chain. Что значит: В отличие от subscribeOn(), имеет значение, куда в цепочку вы помещаете функцию observeOn(), так как этот оператор только изменяет поток, который используется observables, которые следуют ниже. Например, если вы вставляете в свою цепочку следующий код, то каждый observable, который появляется в цепочке с этого момента, будет использовать новый поток. .observeOn(Schedulers.newThread()) Эта цепочка будет продолжать работать в новом потоке, пока не встретится другой оператор observOn(), после чего она переключится на поток, указанный этим оператором. Вы можете управлять потоком, куда конкретные observables отправляют свои уведомления, путём вставки в вашу цепочку нескольких операторов observeOn().

понедельник, 16 декабря 2019 г.

Асинхронное чтение/запись в Realm используя RXJava 2

#java #android #rxjava #realm #rxandroid


Моя первая реализация асинхронной работы с помощью RXJava 2. 

Цель:

Получить json данные с сервера библиотекой Retrofit2. Если успешно, то записать в
Realm и сразу после записи получить обратно данные и отправить адаптеру RecyclerView.

Так вот, я все это реализовал таким образом:

private void fetchChatsFromNetwork(int count, AccessDataModel accessDataModel) {

    String accessToken = accessDataModel.getAccessToken();

    MyApplication.getRestApi().getChats(count, accessToken, Constants.api_version)
            .subscribeOn(Schedulers.io())
            .observeOn(AndroidSchedulers.mainThread())
            .subscribeWith(new DisposableSubscriber() {
                @Override
                public void onNext(ChatsModel chatsModel) {
                    if (chatsRepository.hasData()) {

                        chatsRepository.updateChatsData(chatsModel)
                                .subscribe(new DisposableObserver() {
                                    @Override
                                    public void onNext(ChatsModel localChatsModel) {
                                        Log.d(TAG, "DO, onSuccess updated!");
                                        iGetChatsCallback.onGetChatsSuccess(localChatsModel);
                                    }

                                    @Override
                                    public void onError(Throwable e) {
                                        Log.d(TAG, "DO, onError when update!");
                                        iGetChatsCallback.onGetChatsError(e.getMessage());
                                    }

                                    @Override
                                    public void onComplete() {
                                        dispose();
                                        Log.d(TAG, "DO, onComplete!");
                                    }
                                });

                    } else {
                        chatsRepository.insertChatsData(chatsModel)
                                .subscribe(new DisposableObserver() {
                                    @Override
                                    public void onNext(ChatsModel localChatsModel) {
                                        iGetChatsCallback.onGetChatsSuccess(localChatsModel);
                                        Log.d(TAG, "DO, onSuccess inserted!");
                                    }

                                    @Override
                                    public void onError(Throwable e) {
                                        iGetChatsCallback.onGetChatsError(e.getMessage());
                                        Log.d(TAG, "DO, onError when inserting!");
                                    }

                                    @Override
                                    public void onComplete() {
                                        dispose();
                                        Log.d(TAG, "DO, onComplete!");
                                    }
                                });
                    }
                }

                @Override
                public void onError(Throwable t) {
                    Log.d(TAG, "onError" + t.getMessage());
                }

                @Override
                public void onComplete() {
                    Log.d(TAG, "onComplete");
                }
            });
}


Я записываю данные в Realm в методе onNext() подписчика MyApplication.getRestApi().getChats().

Вот код записи:

public Observable updateChatsData(final ChatsModel chatsModel) {

    return Observable.create(new ObservableOnSubscribe() {
        @Override
        public void subscribe(ObservableEmitter e) throws Exception {
            if (chatsModel != null) {
                realm.executeTransactionAsync(
                        realm -> realm.copyToRealmOrUpdate(chatsModel),
                        () -> {
                            Log.d(LOG_TAG, "Data success updated!");
                            ChatsModel localChatsModel = getAllChatsData();
                            e.onNext(localChatsModel);
                            e.onComplete();
                        },
                        error -> {
                            Log.d(LOG_TAG, "Update data failed!");
                            e.onError(error);
                        });
            }

        }
    });

}


Метод updateChatsData() записывает асинхронно и объявлен в другом классе. 

Как видите мой метод fetchChatsFromNetwork() написан громоздко или мне так кажется.

Вопрос:

Правильно ли я делаю или нет, если нет, то как было бы правильнее?
    


Ответы

Ответ 1



Можно полностью отвязать запись в БД от уведомления адаптера о новых данных. Подпишитесь на Observable, выдающий выборку из БД и уведомляющий о ней адаптер. При сетевом запросе полученные данные пишите в БД. При таком способе Observable из первого пункта уведомит адаптер сразу после записи/обновлении данных в БД. Саму запись в БД макже можно проще сделать через flatMap как-то так: MyApplication.getRestApi().getChats(count, accessToken, Constants.api_version) .flatMap(data -> (chatsRepository.hasData() ? chatsRepository.updateChatsData(data) : chatsRepository.insertChatsData(data)).flatMap(data -> Observable.just(true))) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribeWith( aBoolean -> System.out.println("data in DB updated"), error -> System.out.println("error: " + e.getMessage()) ); Тут, возможно, придётся поиграться с заменой транзакций записей в БД с синхронных на асинхронные (скорее наоборот) из-за того, как работают асинхронные в потоках без Looper.

воскресенье, 15 декабря 2019 г.

Как правильно пересоздать закэшированный Observable используемый вместе с Retrofit?

#java #android #retrofit #rxjava #rxandroid


Дано:

API, возвращающее список с данными в формате JSON.

Задача:

Получить эти данные силами Retrofit+RxJava.

Проблема:

Необходимо сделать изначально один запрос и не дублировать его, если экран будет
повёрнут до окончания задачи. Также нужно иметь возможность перезапустить задачу.  

Что получилось:

Первое я решил, сделав Singlton и закешировав единственный экземпляр Observable с
помощью cache(). 

Второе - полным пересозданием объекта Retrofit (1), экземпляра retrofit-интерфейса
(2) и самого Observable(3). Если 1 и 2 не сделать - 3 - остаётся прежним и возвращает
закэшированные данные, вместо новых.

Вопрос:

Использованный мной способ перезапуска задачи получения данных выглядит плохо. Как
лучше/правильнее пересоздать Observable?



Синглтон для получения/пересоздания Observalbe:

public class SingltonRetrofit
{
    private static RxJavaCallAdapterFactory rxAdapter = RxJavaCallAdapterFactory.createWithScheduler(Schedulers.io());

    private static Gson gson = new GsonBuilder().create();

    private static Retrofit retrofit = new Retrofit.Builder()
            .baseUrl(Const.BASE_URL)
            .addConverterFactory(GsonConverterFactory.create(gson))
            .addCallAdapterFactory(rxAdapter)
            .build();

    private static GetModels apiService = retrofit.create(GetModels.class);
    private static Observable> observableModelsList;

    public static void reset()
    {
        retrofit = new Retrofit.Builder()
                .baseUrl(Const.BASE_URL)
                .addConverterFactory(GsonConverterFactory.create(gson))
                .addCallAdapterFactory(rxAdapter)
                .build();
        apiService = retrofit.create(GetModels.class);
        observableModelsList = null;
    }

    public static Observable> getModelsObservable()
    {
        if (observableModelsList == null)
        {
            observableModelsList = apiService.getModelsList().cache();
        }
        return observableModelsList;
    }
}


P.S.

Этот же вопрос на английском: How to recreate or reset cached Observable, used with
Retrofit to get new data?
    


Ответы

Ответ 1



Не претендую на лучшую реализацию, но... Лично я у себя во избежания повторных запросов держу данные отдельно от любой активити, дабы не захламлять код. Есть некий класс, который всегда создается в памяти при старте приложения к примеру пусть будет Model public class Model { public final MyData mayData; public Model(RestApi mRestApi) { //мне удобно передовать сразу интерфейс апи mayData = new MyData(mRestApi); } } Старт его происходит в классе наследованного от Application. Класс MyData - обычный класс в котором описаны запросы для конкретного экрана или типа данных. public class MyData { public MainPageInfo(RestApi retrofit) { super(retrofit); } public void getData() { mRestApi.getData().enqueue(new Callback>() { @Override public void onResponse(Call> call, Response> response) { } @Override public void onFailure(Call> call, Throwable t) { setError(t); } }); } } Для удобства доступа к модели я храню на нее ссылки в классах унаследованных от фрагмента или активити, в свою очередь их ставлю в наследники для нужных мне вью. Как итог - во время переворота дестроятся сама вьюшка, но не данные и их легко попросить у класса MyData с просто проверкой на null (делать запрос если данных нет). Конечно реализация не совершенна, как и здесь я описал не весь код (кроме всего есть модель слушателей)). Но я почему то думаю для вас это не проблема и вы сами решите подойдет вам этот метод или нет. Единственный недостаток - если дестроется все приложение, то и данные тоже. Зато можно обратиться к данным из любого места =), Если что пишите на почту ) расскажу подробней )

Ответ 2



В итоге сделал так: Как верно написал @mit, метод cache() кэшировал запрос и пересоздание observable в итоге не перезапускало сетевой запрос. При этом в доках и в интернетах я нигде сему упоминания не находил. Видать это как-то связано с внутренней логикой связки Retrofit+OkHttp При этом как верно предложил @Yura Ivanov, потребовался BehaviorSubject. На него подписывается фрагмент и его же можно пересоздать в случае нужды в свежих данных (без пересоздания будет отданы последние данные). При этом при каждом создании/пересоздании BehaviorSubject создаётся Subscriber для получения данных из сети через Observable, создаваемый Retrofit-ом. И он в onError и в onNext вызывает соответствующие методы у BehaviorSubject. При этом не транслируя onComplete, т.к. в этом случае может произойти ситуация, когда данные придут в процессе пересоздания фрагмента и фрагмент получит только последнее событие BehaviorSubject, т.е. событие завершения последовательности, вместо последних полученных данных. Т.е. соединять подпиской напрямую BehaviorSubject и Observable, получающий сетевые данные не стоит. Итого все требования соблюдены: При поворотах экрана будет запущена всего однажды задача на скачивание данных, фрагмент получит данные (или сообщение об ошибке) в любом случае и пересоздавать объекты Retrofita-а не нужно. Итоговый синглтон: public class SingltonRetrofitNew { private static RxJavaCallAdapterFactory rxAdapter = RxJavaCallAdapterFactory.createWithScheduler(Schedulers.io()); private static Gson gson = new GsonBuilder().create(); private static Retrofit retrofit = new Retrofit.Builder() .baseUrl(Const.BASE_URL) .addConverterFactory(GsonConverterFactory.create(gson)) .addCallAdapterFactory(rxAdapter) .build(); private static GetModels apiService = retrofit.create(GetModels.class); private static BehaviorSubject> observableModelsList; private static Observable> observable = apiService.getModelsList(); private static Subscription subscription; private SingltonRetrofitNew() { } public static void resetObservable() { observableModelsList = BehaviorSubject.create(); if (subscription != null && !subscription.isUnsubscribed()) { subscription.unsubscribe(); } subscription = observable.subscribe(new Subscriber>() { @Override public void onCompleted() { //do nothing } @Override public void onError(Throwable e) { observableModelsList.onError(e); } @Override public void onNext(ArrayList hotels) { observableModelsList.onNext(hotels); } }); } public static Observable> getModelsObservable() { if (observableModelsList == null) { resetObservable(); } return observableModelsList; } } Сокращённый фрагмент: public class FragmentsList extends Fragment { private static final String TAG = FragmentList.class.getSimpleName(); private Subscription subscription; private RecyclerView recyclerView; private SwipeRefreshLayout swipeRef; private ArrayList models = new ArrayList<>(); private boolean isLoading; @Nullable @Override public View onCreateView(LayoutInflater inflater, @Nullable ViewGroup container, @Nullable Bundle savedInstanceState) { View v = inflater.inflate(R.layout.fragment, container, false); //init views recyclerView = (RecyclerView) v.findViewById(R.id.recycler); swipeRef = (SwipeRefreshLayout) v.findViewById(R.id.swipe_ref); swipeRefreshLayout.setOnRefreshListener(new OnRefreshListener() { @Override public void onRefresh() { SingltonRetrofitNew.reset(); getModelsList(); } }); if (savedInstanceState != null) { models = savedInstanceState.getParcelableArrayList(Const.KEY_MODELS); isLoading = savedInstanceState.getBoolean(Const.KEY_IS_LOADING); } if (models.size() == 0 || isLoading) { getModelsList(); } //TODO show saved data if is return v; } @Override public void onDestroy() { super.onDestroy(); if (subscription != null && !subscription.isUnsubscribed()) { subscription.unsubscribe(); } } private void getModelsList() { isLoading = true; swipeRef.setRefreshing(true); if (subscription != null && !subscription.isUnsubscribed()) { subscription.unsubscribe(); } subscription = SingltonRetrofitNew.getModelsObservable(). subscribeOn(Schedulers.io()). observeOn(AndroidSchedulers.mainThread()). subscribe(new Subscriber>() { @Override public void onCompleted() { Log.d(TAG, "onCompleted"); } @Override public void onError(Throwable e) { Log.d(TAG, "onError", e); isLoading = false; swipeRef.setRefreshing(false); Snackbar.make(recyclerView, R.string.connection_error, Snackbar.LENGTH_SHORT) .setAction(R.string.try_again, new View.OnClickListener() { @Override public void onClick(View v) { SingltonRetrofitNew.reset(); getModelsList(); } }) .show(); } @Override public void onNext(ArrayList newModels) { isLoading = false; swipeRef.setRefreshing(false); models.clear(); models.addAll(newModels); //TODO show data } }); } @Override public void onSaveInstanceState(Bundle outState) { super.onSaveInstanceState(outState); outState.putParcelableArrayList(Const.KEY_MODELS, models); outState.putBoolean(Const.KEY_IS_LOADING, isLoading); } } Всё вместе на gitHub: RxRetrofitAndScreenOrientation Статья на ХабраХабр про решение: Используем RxJava и Retrofit на Android, учитывая поворот экрана

пятница, 13 декабря 2019 г.

Назначение RxJava, RxAndroid

#java #android #rxjava #rxandroid


Для чего нужны библиотеки RxJava и RxAndroid? Где их применяют?
    


Ответы

Ответ 1



RxJava - либа для обработки событий как последовательностей. Надстройка над паттерном "наблюдатель", умеет по-всякому преобразовывать последовательности событий и так далее. Либу можно заюзать где угодно, но особо полезна когда нужно делать всякие сложные преобразования над событиями, типа смешать, сгруппировать, фильтровать, ограничивать и т.п. RxAndroid - либа на базе RxJava, содержит всякие штуки для андроида типа шедулера для гуевого потока. Всякая логика типа "получать символы в текстбоксах от юзера, отправлять запросы к серверу в другом потоке, но не слишком часто, с какой-то задержкой, обновлять интерфейс" пишутся на этих либах в пару строк, без гемороя с кучей служебного кода.

Ответ 2



RxJava это набор инструментов, которые помогут написать красивый, но более простой код. Он обладает интуитивностью, расширяемостью, декларативностью, возможностью компоновки и конвертируемостью. На эту тему нашла полезную статью, которая как раз раскрывает назначение этой Java библиотеки.

среда, 27 ноября 2019 г.

Как реализовать очередь запросов?

#java #android #rxjava #rxandroid


Исходные данные:

Есть проект, который использует методы из SDK стороннего проекта (асинхронные методы
для работы с их сервером с параметром-колбэком) и содержит в себе Retrofit2 + OkHttp
+ Rx для работы с сервером напрямую. Для удобства используются лямбды.

Задача:

Необходимо, чтобы все запросы (и из SDK, и через Retrofit) выполнялись не чаще чем
5 раз в секунду, а если лимит превышен, выполнялись с задержкой.

Вопрос: 

Как это реализовать? Первое, что приходит в голову - это Service + BroadcastReceiver.
Но придется слушать ресивер в каждом активити/фрагменте, плюс получится не такая удобная
реализация колбэков... Может быть можно это как-то реализовать с помощью Rx (в том
числе обернув методы SDK), чтобы оставить лямбды для удобства?
    


Ответы

Ответ 1



Может это поможет? Subscription subscription = Observable.interval(1000, 5000, TimeUnit.MILLISECONDS) .observeOn(AndroidSchedulers.mainThread()) .subscribe(new Action1() { public void call(Long aLong) { // here is the task that should repeat } }); https://stackoverflow.com/a/49718536/10965132

пятница, 12 июля 2019 г.

Как правильно делать запросы в операторе flatMap?

Я получаю с сервера список чатов(переписок). Их запихиваю в класс ChatsModel. И в этом классе имеется userId, то есть в каждом чате имеется id того пользователя. Других информации о пользователе нет.И еще есть одно поле для данных о пользователе типа UserMessagesResponse
Так вот для начала этот userMessagesResponse равен null
После того как я получил список чатов и userId, я должен отправить другой запрос для получения информации о пользователе. Делаю я это таким образом:
private void loadChatsFromNetwork(int count, AccessDataModel accessDataModel) { String accessToken = accessDataModel.getAccessToken();
Flowable chatsModelSingle = getChatsApi().getChats(count, accessToken, Constants.api_version) .subscribeOn(Schedulers.io()) .flatMap(chatsModel -> { RealmList items = chatsModel.getResponse().getItems(); StringBuilder userIds = new StringBuilder();
for (Item item : items) { userIds.append(item.getMessage().getUserId()).append(","); }
return loadUsersById(userIds, chatsModel); }) .observeOn(AndroidSchedulers.mainThread());
chatsModelSingle.subscribe(chatsModel -> { Log.d(TAG, chatsModel.getResponse().getItems().first().getMessage().getMessagesUserItem().getFirstName()); chatsRepository.updateChatsData(chatsModel); iGetChatsCallback.onGetChatsSuccess(chatsModel); }, throwable -> { iGetChatsCallback.onGetChatsError(throwable.getMessage()); Log.d(TAG, "onError() " + throwable.getMessage()); }); }
private ChatsModel loadUsersById(StringBuilder userIds, ChatsModel chatsModel) {
MyApplication.getChatsApi().getUsersByChats(userIds.toString(), "photo_100") .subscribe(messagesUser -> {
RealmList item = chatsModel.getResponse().getItems();
for (int i = 0; i < item.size(); i++) { Message message = item.get(i).getMessage();
RealmList messagesUserItemList = messagesUser.getUserMessagesResponse(); for (UserMessagesResponse messagesResponse : messagesUserItemList) { if (messagesResponse.getUid().equals(message.getUserId())) { message.setMessagesUserItem(messagesResponse); chatsModel.getResponse().getItems().get(i).setMessage(message); } } }
});
return chatsModel; }
Все эти действия происходят в операторе flapMap, так как мне нужно полученную информацию о пользователе запихнуть в поле userMessagesResponse класса ChatsModel. И в случае успеха отправляю в adapter.
Оба запроса корректно работают. Получаю список userid, получаю данные о пользователе.
Проблема в том, после возвращения chatsModel в flatMap, где return chatsModel, дальше ничего не происходит, то есть до подписчика не доходит ничего, точнее подписчик никак не реагирует.
Вопрос: Как исправить это и вообще как правильно решать такого рода задачи?


Ответ

Если формально, то проблема ваша в том, что вы делаете return chatsModel до того момента, как выполнится код из subscribe, в котором он наполняется.
Если по сути, то ошибка у вас в идеологии. Вы пытаетесь скрестить ежа с ужом. Идеология ReactiveX состоит в том, что вы управляете не фиксированными объектами, а потоками данных(событий), причем асинхронно. А у вас получается так: пошел асинхронный поток, вы его тут-же пытаетесь собрать в объект, синхронно причем, и запустить в другой поток.
Распишем ваш поток:
Происходит какое-то событие, которое инициирует загрузку чатов Получаем список чатов List Для каждого чата получаем информацию о пользователе (getUserInfo) Делаем с этими чатами какую-то полезную работу (redrawMyNiceChatsTable)
Т.е. логика у вас должны быть примерно следующая(отчасти псевдокод):
class MyMainClass{ private ChatsUpdater updater = new ChatsUpdater(); private List chats = new ArrayList<>();
private void onSomeEventOccured(){ chats.clear(); updater.startUpdate() .subscribe( {chat -> chats.add(chat)}, // onNext {redrawMyNiceChatsTable(chats)} // onComplete ) } }
class ChatsUpdater{ public Flowable startUpdate(){ return getChats().flatMap( chat -> getUserInfo(chat), (chat, userInfo) -> { chat.setUserInfo(userInfo); return Observable.just(chat); } ) }
private Flowable getChats(){ ..... }
private Flowable getUserInfo(Chat chat){ ..... } }
UPD Ну а если совсем по уму, то, так как у вас жесткая, обязательная, цепочка Chat->UserInfo, вам надо переделать API на стороне сервера, чтобы оно отдавало сразу всю необходимую информацию.

воскресенье, 7 июля 2019 г.

Как выполнить метод после нескольких запросов?

Хотел бы реализовать: По нажатию кнопки происходит загрузка информации с API и добавление (Если пустая)/обновление (Если заполненная) этой информации в базу данных.
Более подробно: Использую Retrofit для запросов и БД на SQLLite. Насколько я понял, нужно использовать фоновый асинхронный поток после нажатия кнопки и там загружать данные и заполнять\обновлять БД. Много читал похожих вопросов, везде советуют RxJava, сидел разбирался в RxJava, принцип понял, как он устроен, но примера для нескольких запросов не нашел.
Думал, сначала что все будет обновляться в UI потоке (приложение все равно не имеет практической ценности без БД), а потом подумал, что пусть оно там само обновляется в фоне, а потом какой-нибудь Toast вылезет, мол удачно все прошло. Так логичнее.
В общем, вот код образно:
api.getTransportTypes(JSON).enqueue(new Callback>() {...}; api.getMarshes(JSON).enqueue(new Callback>() {...}; api.getStops(JSON).enqueue(new Callback>() {...}; ...
Log.d("MainActivity", "Data from Api downloaded.");
UpdateDB();
Log.d("MainActivity", "Data Base updated.");
Вопрос: Как мне объединить и выполнить все запросы в API, а после, что бы вызвался метод, по окончанию загрузки, например образный - UpdateDB();?
Если можно какой-нибудь актуальный пример? С сегодняшними фреймворками? На RxJava или может какой-то способ через AsyncTask, Handler, Loader? Или алгоритм, хотя бы подробный (Хотя без кода все равно не понятно будет, наверное).


Ответ

Заюзать rx, как вариант, можно так:
1) Реализовать запросы через RxJava
Добавить .addCallAdapterFactory(RxJava2CallAdapterFactory.create()) к Retrofit.Builder() Поменять методы в сервис интерфейсе @GET("url") fun getStops(): Flowable>
2) Применить zip к полученным flowable
Flowable.zip(api.getMarshes(), api.getStops(), api.getTransportTypes(), Function3, List, List, Triple, List, List>> { marshes, stops, transportTypes -> Triple(marshes, stops, transportTypes) }) .flatMapCompletable { insertToDb(it) }
Zip отправит запросы параллельно, дождется всех ответов и заэмитит дальше
3) Записать в базу
fun insertToDb(triple: Triple, List, List>): Completable { return Completable.fromCallable { database.insert(triple.first) database.insert(triple.second) database.insert(triple.third) } }
Если ваша БД поддерживает rx, то Flowable.fromCallable можно избежать PS код на kotlin

воскресенье, 2 июня 2019 г.

AlertDialog с помощью RxJava

С помощью AsyncTask при обработке массива данных мы легко можем вывести с помощью AlertDialog прогресс выполнения обработки.
Можно ли это сделать с помощью RxJava? Прошу вашего простейшего примера.
И как решается проблема поворота девайса при загрузке данных с сервера при помощи RxJava и отображения прогресса на AlertDialog?


Ответ

Самый простой способ это заюзать метод from:
List list = new ArrayList<>(); for (int i = 0; i < 100; i++) { list.add("item "+i); }
Observable.from(list) .observeOn(AndroidSchedulers.mainThread()) .subscribe(new Subscriber() { @Override public void onCompleted() { if(progress!=null && progress.isShowing()) progress.dismiss(); }
@Override public void onError(Throwable e) { if(progress!=null && progress.isShowing()) progress.dismiss(); }
@Override public void onNext(String s) { progress = ProgressDialog.show(this, "dialog title", "dialog message", true); } });

понедельник, 27 мая 2019 г.

java.lang.IllegalStateException: Expected BEGIN_ARRAY but was BEGIN_OBJECT (Retrofit 2)

Получаю такую ошибку java.lang.IllegalStateException: Expected BEGIN_ARRAY but was BEGIN_OBJECT, когда получаю данные с сервера. Чем вызвана эта ошибка? И как ее исправить?
JSON, который приходит с сервера:
{ "group": [ { "name": "Group1", "description": "Test Group 1" }, { "name": "Group2", "description": "Group Name Updated" } ] }
Интерфейс API:
@GET("groups") Observable > getAllGroups(@Header("Authorization") String auth, @Header("Content-type") String contentType, @Header("Accept") String accept );
Метод, в котором получаю данные:
private void getAllGroups() { String group = "Group1"; String credentials = "admin" + ":" + "admin"; final String basic = "Basic " + Base64.encodeToString(credentials.getBytes(), Base64.NO_WRAP); String contentType = "application/json"; String accept = "application/json"; Subscription subscription = App.service.getAllGroups(basic, contentType, accept) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .subscribe(groups -> { groupList.addAll(groups); }, throwable -> { Log.e("All group error", String.valueOf(throwable)); });
addSubscription(subscription); }
Класс Group:
public class Group {
private String name; private String description; private String admins; private Member members;
public Group(String name, String description, String admins, Member members) { this.name = name; this.description = description; this.admins = admins; this.members = members; }
public Group(String name, String description) { this.name = name; this.description = description; } // getters setters }


Ответ

Сервер присылает JSONObject, в котором JSONArray, а вы пытаетесь парсить как JSONArray. Должно быть как то так:
public class GroupList { ArrayList group; //setters, getters }

public class Group {
private String name; private String description;
... }

пятница, 17 мая 2019 г.

Как превратить EditText в поток данных для RxJava2?

Начал постепенно переходить на реактивное программирование с замечательным фреймворком RxJava 2 и заинтересовал простейший пример создания потока в динамическом стиле. К сожалению, не считаю пример Observable.just(1,2,3) достаточно полным для себя, но вот уже несколько часов мучаю себя мыслью, как превратить изменения в EditText в поток данных? Если не сложно, опишите, пожалуйста, пример как сделать реальный Observable из TextWatcher, чтобы можно было на него подписаться и подписчик реагировал на изменения текста в EditText.


Ответ

Можно воспользоваться готовой библиотекой RxBinding. Также в ней можно посмотреть конкретную реализацию
Вообще для связки callback-методов с rx используется метод Observable.create, например:
Observable.create(emitter -> { TextWatcher watcher = new TextWatcher() { @Override public void beforeTextChanged(CharSequence charSequence, int i, int i1, int i2) { }
@Override public void onTextChanged(CharSequence charSequence, int i, int i1, int i2) { }
@Override public void afterTextChanged(Editable editable) { if (!emitter.isDisposed()) { //если еще не отписались emitter.onNext(editable.toString()); //отправляем текущее состояние } } }; emitter.setCancellable(() -> editText.removeTextChangedListener(watcher)); //удаляем листенер при отписке от observable editText.addTextChangedListener(watcher); });
Нужно не забыть подписаться, и самое главное, отписаться от такого источника данных, т.к. он держит ссылку на editText и может привести к утечке памяти.