From 049ebd83baab4be801713d1fdc7d5fe5f10999c1 Mon Sep 17 00:00:00 2001 From: AlxTsaregorodtsev <101584665+AlxTsaregorodtsev@users.noreply.github.com> Date: Wed, 11 May 2022 23:07:31 +0700 Subject: [PATCH 1/2] Create README.md --- README.md | 2 ++ 1 file changed, 2 insertions(+) create mode 100644 README.md diff --git a/README.md b/README.md new file mode 100644 index 0000000..0b526bc --- /dev/null +++ b/README.md @@ -0,0 +1,2 @@ +# RxJavaHomework +# From ed7cd41ba890a33ece06a338f78e7d4c68148ea2 Mon Sep 17 00:00:00 2001 From: AlxTsaregorodtsev <101584665+AlxTsaregorodtsev@users.noreply.github.com> Date: Wed, 11 May 2022 23:38:00 +0700 Subject: [PATCH 2/2] Add files via upload --- Main.kt | 159 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 159 insertions(+) create mode 100644 Main.kt diff --git a/Main.kt b/Main.kt new file mode 100644 index 0000000..c88466e --- /dev/null +++ b/Main.kt @@ -0,0 +1,159 @@ +import io.reactivex.rxjava3.core.* +import java.rmi.server.ServerNotActiveException +import java.util.* +import java.util.concurrent.TimeUnit +import kotlin.random.Random + +private const val LATENCY = 700L +private const val RESPONSE_LENGTH = 2048 + +fun main() { + // Функции можно вызывать отсюда для проверки + // для ДЗ лучше использовать blockingSubscribe вместо subscribe потому что subscribe подпишется на изменения, + // но изменения в большинстве случаев будут получены позже, чем выполнится функция main, поэтому в консоли ничего + // не будет выведено. blockingSubscribe работает синхронно, поэтому результат будет выведен в консоль + // + // В реальных программах нужно использовать subscribe или передавать данные от источника к источнику для + // асинхронного выполнения кода. + // + // Несмотря на то, что в некоторых заданиях фигурируют слова "синхронный" и "асинхронный" в рамках текущего ДЗ + // это всего лишь имитация, реальное переключение между потоками будет рассмотрено на следующем семинаре + + println("#1") + requestDataFromServerAsync() + .blockingSubscribe( + { + println(it) + }, + { + println(it.message) + } + ) + println("--------------") + println("#2") + requestServerAsync() + .blockingSubscribe( + { + println("OK") + }, + { + println(it.message) + } + ) + println("--------------") + println("#3") + requestDataFromDbAsync() + .blockingSubscribe( + { + println("OK") + println(it) + }, + { + println(it.message) + }, + { + println("complete") + } + ) + println("--------------") + println("#4") + emitEachSecond() + + + xMap { flatMapCompletable(it) } + xMap { concatMapCompletable (it) } + xMap { switchMapCompletable(it) } +} + +// 1) Какой источник лучше всего подойдёт для запроса на сервер, который возвращает результат? +// Почему? Single - исполнится и завершится событием Success или Error +// Дописать функцию +fun requestDataFromServerAsync(): Single { + + // Функция имитирует синхронный запрос на сервер, возвращающий результат + fun getDataFromServerSync(): ByteArray? { + Thread.sleep(LATENCY) + val success = Random.nextBoolean() + return if (success) Random.nextBytes(RESPONSE_LENGTH) else null + } + + + return Single.fromCallable { getDataFromServerSync() } + +} + + +// 2) Какой источник лучше всего подойдёт для запроса на сервер, который НЕ возвращает результат? +// Почему? Completable - важено только выплнение запроса а не результат +// Дописать функцию +fun requestServerAsync(): Completable { + + // Функция имитирует синхронный запрос на сервер, не возвращающий результат + fun getDataFromServerSync() { + Thread.sleep(LATENCY) + if (Random.nextBoolean()) throw ServerNotActiveException() + } + return Completable.fromAction { getDataFromServerSync() } + +} + +// 3) Какой источник лучше всего подойдёт для однократного асинхронного возвращения значения из базы данных? +// Почему? Maybe - из БД можно получить null +// Дописать функцию +fun requestDataFromDbAsync(): Maybe /* -> ??? */ { + + // Функция имитирует синхронный запрос к БД не возвращающий результата + fun getDataFromDbSync(): T? { + Thread.sleep(LATENCY); return null + } + + return Maybe.fromCallable { getDataFromDbSync() } + +} + +// 4) Примените к источнику оператор (несколько операторов), которые приведут к тому, чтобы элемент из источника +// отправлялся раз в секунду (оригинальный источник делает это в 2 раза чаще). +// Значения должны идти последовательно (0, 1, 2 ...) +// Для проверки результата можно использовать .blockingSubscribe(::printer) +fun emitEachSecond() { + + // Источник + fun source(): Flowable = Flowable.interval(500, TimeUnit.MILLISECONDS) + + // Принтер + fun printer(value: Long) = println("${Date()}: value = $value") + + source().filter { it % 2 == 0L } + .map { it / 2 } + .blockingSubscribe(::printer) +} + +// 5) Функция для изучения разницы между операторами concatMap, flatMap, switchMap +// Нужно вызвать их последовательно и разобраться чем они отличаются +// Документацию в IDEA можно вызвать встав на функцию (например switchMap) курсором и нажав hotkey для вашей раскладки +// Mac: Intellij Idea -> Preferences -> Keymap -> Быстрый поиск "Quick documentation" +// Win, Linux: File -> Settings -> Keymap -> Быстрый поиск "Quick documentation" +// +// конструкция в аргументах функции xMap не имеет значения для ДЗ и создана для удобства вызова функции, чтобы была +// возможность удобно заменять тип маппинга +// +// Вызов осуществлять поочерёдно из функции main +// +// xMap { flatMapCompletable(it) } +// xMap { concatMapCompletable (it) } +// xMap { switchMapCompletable(it) } +// +fun xMap(mapper: Flowable.(internalMapper: (Int) -> Completable) -> Completable) { + + fun waitOneSecond() = Completable.timer(1, TimeUnit.SECONDS) + + println("${Date()}: start") + Flowable.fromIterable(0..20) + .mapper { iterableIndex -> + + waitOneSecond() + .doOnComplete { println("${Date()}: finished operation for iterable index $iterableIndex") } + + } + .blockingSubscribe() +} \ No newline at end of file