diff --git a/src/main/kotlin/Main.kt b/src/main/kotlin/Main.kt index 18c0de6..d5ce923 100644 --- a/src/main/kotlin/Main.kt +++ b/src/main/kotlin/Main.kt @@ -1,6 +1,7 @@ import io.reactivex.rxjava3.core.Completable import io.reactivex.rxjava3.core.Flowable import io.reactivex.rxjava3.core.Maybe +import io.reactivex.rxjava3.core.Single import java.rmi.server.ServerNotActiveException import java.util.* import java.util.concurrent.TimeUnit @@ -10,6 +11,13 @@ private const val LATENCY = 700L private const val RESPONSE_LENGTH = 2048 fun main() { + //requestDataFromServerAsync() + //requestServerAsync() + //requestDataFromDbAsync() + //emitEachSecond() + //xMap { flatMapCompletable(it) } + //xMap { concatMapCompletable(it) } + //xMap { switchMapCompletable(it) } // Функции можно вызывать отсюда для проверки // для ДЗ лучше использовать blockingSubscribe вместо subscribe потому что subscribe подпишется на изменения, // но изменения в большинстве случаев будут получены позже, чем выполнится функция main, поэтому в консоли ничего @@ -23,46 +31,71 @@ fun main() { } // 1) Какой источник лучше всего подойдёт для запроса на сервер, который возвращает результат? -// Почему? +// Почему? Single сгенерирует ошибку, в случае отсутсвия значение, а мы всегда ожидаем результат // Дописать функцию -fun requestDataFromServerAsync() /* -> ??? */ { +fun requestDataFromServerAsync(): ByteArray/* -> ??? */ { // Функция имитирует синхронный запрос на сервер, возвращающий результат fun getDataFromServerSync(): ByteArray? { - Thread.sleep(LATENCY); + Thread.sleep(LATENCY) val success = Random.nextBoolean() return if (success) Random.nextBytes(RESPONSE_LENGTH) else null } + lateinit var result: ByteArray + Single.fromCallable { getDataFromServerSync() } + .blockingSubscribe { + println(it) + result = it + } /* return ??? */ + return result } // 2) Какой источник лучше всего подойдёт для запроса на сервер, который НЕ возвращает результат? -// Почему? +// Почему? Completable потому что нас интересует лишь факт выполнения операции // Дописать функцию -fun requestServerAsync() /* -> ??? */ { - +fun requestServerAsync(): Boolean/* -> ??? */ { // Функция имитирует синхронный запрос на сервер, не возвращающий результат fun getDataFromServerSync() { Thread.sleep(LATENCY) if (Random.nextBoolean()) throw ServerNotActiveException() } - /* return ??? */ + var result = true + Completable.fromAction { getDataFromServerSync() } + .blockingSubscribe({ + println("requestServerAsync success") + }, { + result = false + it.printStackTrace() + }) + return result } // 3) Какой источник лучше всего подойдёт для однократного асинхронного возвращения значения из базы данных? -// Почему? +// Почему? Maybe мы ожидаем возможное null значение // Дописать функцию -fun requestDataFromDbAsync() /* -> ??? */ { - +fun requestDataFromDbAsync(): T? /* -> ??? */ { // Функция имитирует синхронный запрос к БД не возвращающий результата fun getDataFromDbSync(): T? { Thread.sleep(LATENCY); return null } + var result: T? = null + Maybe.fromCallable { getDataFromDbSync() } + .blockingSubscribe({ + println("requestDataFromDbAsync $it success") + result = it as T + }, + { println("requestDataFromDbAsync error") }, + { + println("requestDataFromDbAsync complete") + result = null + }) /* return */ + return result } // 4) Примените к источнику оператор (несколько операторов), которые приведут к тому, чтобы элемент из источника @@ -78,6 +111,10 @@ fun emitEachSecond() { fun printer(value: Long) = println("${Date()}: value = $value") // code here + source() + .onBackpressureBuffer() + .concatMap { Flowable.interval(1000, TimeUnit.MILLISECONDS) } + .blockingSubscribe({ printer(it) }, { it.printStackTrace() }) } // 5) Функция для изучения разницы между операторами concatMap, flatMap, switchMap @@ -102,7 +139,6 @@ fun xMap(mapper: Flowable.(internalMapper: (Int) -> Completable) -> Complet println("${Date()}: start") Flowable.fromIterable(0..20) .mapper { iterableIndex -> - waitOneSecond() .doOnComplete { println("${Date()}: finished operation for iterable index $iterableIndex") }