Kotlin работает на Android

В сопрограммах поток — это тип данных, который может последовательно выдавать несколько значений, в отличие от функций приостановки , которые возвращают только одно значение. Например, поток можно использовать для получения обновлений в реальном времени из базы данных.

Потоки строятся на основе сопрограмм и могут передавать несколько значений. Поток — это, по сути, поток данных , который может обрабатываться асинхронно. Выдаваемые значения должны быть одного типа. Например, Flow<Int> — это поток, который выдает целочисленные значения.

Поток очень похож на Iterator , который генерирует последовательность значений, но использует функции приостановки для асинхронной генерации и потребления значений. Это означает, например, что поток может безопасно отправлять сетевой запрос для получения следующего значения, не блокируя основной поток.

В потоке данных участвуют три сущности:

  • Производитель создает данные, которые добавляются в поток. Благодаря сопрограммам, потоки могут также создавать данные асинхронно.
  • (Необязательно) Посредники могут изменять каждое значение, передаваемое в поток, или сам поток.
  • Потребитель получает значения из потока.

Субъекты, участвующие в потоках данных: потребитель, необязательные посредники и производитель.
Рисунок 1. Субъекты, участвующие в потоках данных: потребитель, посредники (по желанию) и производитель.

In Android, a repository is typically a producer of UI data that has the user interface (UI) as the consumer that ultimately displays the data. Other times, the UI layer is a producer of user input events and other layers of the hierarchy consume them. Layers in between the producer and consumer usually act as intermediaries that modify the stream of data to adjust it to the requirements of the following layer.

Создание потока

Для создания потоков используйте API конструктора потоков . Функция конструктора flow создает новый поток, в котором вы можете вручную передавать новые значения в поток данных с помощью функции emit .

В следующем примере источник данных автоматически получает последние новости через фиксированный интервал. Поскольку функция приостановки не может возвращать несколько последовательных значений, источник данных создает и возвращает поток для выполнения этого требования. В данном случае источник данных выступает в роли производителя.

class NewsRemoteDataSource(
    private val newsApi: NewsApi,
    private val refreshIntervalMs: Long = 5000
) {
    val latestNews: Flow<List<ArticleHeadline>> = flow {
        while (true) {
            val latestNews = newsApi.fetchLatestNews()
            emit(latestNews) // Emits the result of the request to the flow
            delay(refreshIntervalMs) // Suspends the coroutine for some time
        }
    }
}

// Interface that provides a way to make network requests with suspend functions
interface NewsApi {
    suspend fun fetchLatestNews(): List<ArticleHeadline>
}

Конструктор flow выполняется внутри сопрограммы. Таким образом, он использует те же асинхронные API, но с некоторыми ограничениями:

  • Потоки данных являются последовательными . Поскольку производитель находится в сопрограмме, при вызове функции приостановки он приостанавливает выполнение до тех пор, пока эта функция не вернет результат. В примере производитель приостанавливает выполнение до завершения сетевого запроса fetchLatestNews . Только после этого результат передается в поток.
  • При использовании конструктора flow производитель не может emit значения из другого CoroutineContext . Поэтому не следует вызывать метод emit в другом CoroutineContext , создавая новые корутины или используя блоки кода withContext . В таких случаях можно использовать другие конструкторы потоков, например callbackFlow .

Модификация потока

Промежуточные операторы могут использовать промежуточные операторы для изменения потока данных без потребления значений. Эти операторы представляют собой функции, которые при применении к потоку данных запускают цепочку операций, которые не выполняются до тех пор, пока значения не будут потреблены в будущем. Подробнее о промежуточных операторах см. в справочной документации Flow .

В приведенном ниже примере слой репозитория использует промежуточную map операторов для преобразования данных, отображаемых в View :

class NewsRepository(
    private val newsRemoteDataSource: NewsRemoteDataSource,
    private val userData: UserData
) {
    /**
     * Returns the favorite latest news applying transformations on the flow.
     * These operations are lazy and don't trigger the flow. They just transform
     * the current value emitted by the flow at that point in time.
     */
    val favoriteLatestNews: Flow<List<ArticleHeadline>> =
        newsRemoteDataSource.latestNews
            // Intermediate operation to filter the list of favorite topics
            .map { news -> news.filter { userData.isFavoriteTopic(it) } }
            // Intermediate operation to save the latest news in the cache
            .onEach { news -> saveInCache(news) }
}

Промежуточные операторы могут применяться один за другим, образуя цепочку операций, которые выполняются лениво при появлении элемента в потоке. Следует отметить, что простое применение промежуточного оператора к потоку не запускает сборку мусора.

Сбор из потока

Используйте терминальный оператор , чтобы запустить поток и начать прослушивание значений. Чтобы получать все значения в потоке по мере их поступления, используйте collect . Подробнее об терминальных операторах можно узнать в официальной документации Flow .

Поскольку collect — это функция с приостановкой выполнения, её необходимо выполнять внутри сопрограммы. Она принимает в качестве параметра лямбда-функцию, которая вызывается при каждом новом значении. Так как это функция с приостановкой выполнения, сопрограмма, вызывающая collect может приостанавливать выполнение до тех пор, пока поток не будет завершен.

Продолжая предыдущий пример, вот простая реализация ViewModel , использующей данные из слоя репозитория:

class LatestNewsViewModel(
    private val newsRepository: NewsRepository
) : ViewModel() {

    init {
        viewModelScope.launch {
            // Trigger the flow and consume its elements using collect
            newsRepository.favoriteLatestNews.collect { favoriteNews ->
                // Update UI with the latest favorite news
            }
        }
    }
}

Сбор потока данных запускает генератор, который обновляет последние новости и выдает результат сетевого запроса с фиксированным интервалом. Поскольку генератор всегда остается активным в цикле while(true) , поток данных будет закрыт, когда ViewModel будет очищен и viewModelScope будет отменен.

Сбор проб воды может прекратиться по следующим причинам:

  • Сопрограмма, выполняющая сбор данных, отменяется, как показано в предыдущем примере. Это также останавливает базовый производитель.
  • Производитель завершает отправку элементов. В этом случае поток данных закрывается, и сопрограмма, вызвавшая collect возобновляет выполнение.

Flows are cold and lazy unless specified with other intermediate operators. This means that the producer code is executed each time a terminal operator is called on the flow. In the previous example, having multiple flow collectors causes the data source to fetch the latest news multiple times on different fixed intervals. To optimize and share a flow when multiple consumers collect at the same time, use the shareIn operator.

Обнаружение непредвиденных исключений

Реализация производителя может быть взята из сторонней библиотеки. Это означает, что она может генерировать неожиданные исключения. Для обработки таких исключений используйте промежуточный оператор catch .

class LatestNewsViewModel(
    private val newsRepository: NewsRepository
) : ViewModel() {

    init {
        viewModelScope.launch {
            newsRepository.favoriteLatestNews
                // Intermediate catch operator. If an exception is thrown,
                // catch and update the UI
                .catch { exception -> notifyError(exception) }
                .collect { favoriteNews ->
                    // Update UI with the latest favorite news
                }
        }
    }
}

В предыдущем примере при возникновении исключения лямбда-функция collect не вызывается, поскольку новый элемент еще не получен.

catch также может emit элементы в поток. В примере с репозиторием вместо этого можно emit кэшированные значения:

class NewsRepository(
    // ...
) {
    val favoriteLatestNews: Flow<List<ArticleHeadline>> =
        newsRemoteDataSource.latestNews
            .map { news -> news.filter { userData.isFavoriteTopic(it) } }
            .onEach { news -> saveInCache(news) }
            // If an error happens, emit the last cached values
            .catch { exception -> emit(lastCachedNews()) }
}

В этом примере при возникновении исключения вызывается лямбда-функция ` collect , поскольку из-за исключения в поток был отправлен новый элемент.

Выполнение в другом контексте сопрограммы

By default, the producer of a flow builder executes in the CoroutineContext of the coroutine that collects from it, and as previously mentioned, it cannot emit values from a different CoroutineContext . This behavior might be undesirable in some cases. For instance, in the examples used throughout this topic, the repository layer shouldn't be performing operations on Dispatchers.Main that is used by viewModelScope .

To change the CoroutineContext of a flow, use the intermediate operator flowOn . flowOn changes the CoroutineContext of the upstream flow , meaning the producer and any intermediate operators applied before (or above) flowOn . The downstream flow (the intermediate operators after flowOn along with the consumer) is not affected and executes on the CoroutineContext used to collect from the flow. If there are multiple flowOn operators, each one changes the upstream from its current location.

class NewsRepository(
    private val newsRemoteDataSource: NewsRemoteDataSource,
    private val userData: UserData,
    private val defaultDispatcher: CoroutineDispatcher
) {
    val favoriteLatestNews: Flow<List<ArticleHeadline>> =
        newsRemoteDataSource.latestNews
            .map { news -> // Executes on the default dispatcher
                news.filter { userData.isFavoriteTopic(it) }
            }
            .onEach { news -> // Executes on the default dispatcher
                saveInCache(news)
            }
            // flowOn affects the upstream flow ↑
            .flowOn(defaultDispatcher)
            // the downstream flow ↓ is not affected
            .catch { exception -> // Executes in the consumer's context
                emit(lastCachedNews())
            }
}

В этом коде операторы onEach и map используют defaultDispatcher , тогда как оператор catch и consumer выполняются в Dispatchers.Main , используемом viewModelScope .

Поскольку слой источника данных выполняет операции ввода-вывода, следует использовать диспетчер, оптимизированный для таких операций:

class NewsRemoteDataSource(
    // ...
    private val ioDispatcher: CoroutineDispatcher
) {

    val latestNews: Flow<List<ArticleHeadline>> = flow {
        // Executes on the IO dispatcher
        // ...
    }
        .flowOn(ioDispatcher)
}

Потоки в библиотеках Jetpack

Flow интегрирован во многие библиотеки Jetpack и популярен среди сторонних библиотек для Android. Flow отлично подходит для обновления данных в реальном времени и обработки бесконечных потоков данных.

С помощью Flow и Room можно получать уведомления об изменениях в базе данных. При использовании объектов доступа к данным (DAO) возвращайте тип Flow для получения обновлений в режиме реального времени.

@Dao
abstract class ExampleDao {
    @Query("SELECT * FROM Example")
    abstract fun getExamples(): Flow<List<Example>>
}

При каждом изменении в таблице Example генерируется новый список с новыми элементами базы данных.

Преобразуйте API на основе обратных вызовов в потоки обработки событий.

callbackFlow — это конструктор потоков, позволяющий преобразовывать API, использующие коллбэки, в потоки. В качестве примера можно привести API Firebase Firestore для Android, использующие коллбэки.

Для преобразования этих API в потоки и отслеживания обновлений базы данных Firestore можно использовать следующий код:

class FirestoreUserEventsDataSource(
    private val firestore: FirebaseFirestore
) {
    // Method to get user events from the Firestore database
    fun getUserEvents(): Flow<UserEvents> = callbackFlow {

        // Reference to use in Firestore
        var eventsCollection: CollectionReference? = null
        try {
            eventsCollection = FirebaseFirestore.getInstance()
                .collection("collection")
                .document("app")
        } catch (e: Throwable) {
            // If Firebase cannot be initialized, close the stream of data
            // flow consumers will stop collecting and the coroutine will resume
            close(e)
        }

        // Registers callback to firestore, which will be called on new events
        val subscription = eventsCollection?.addSnapshotListener { snapshot, _ ->
            if (snapshot == null) { return@addSnapshotListener }
            // Sends events to the flow! Consumers will get the new events
            try {
                trySend(snapshot.getEvents())
            } catch (e: Throwable) {
                // Event couldn't be sent to the flow
            }
        }

        // The callback inside awaitClose will be executed when the flow is
        // either closed or cancelled.
        // In this case, remove the callback from Firestore
        awaitClose { subscription?.remove() }
    }
}

В отличие от конструктора flow , callbackFlow позволяет передавать значения из другого CoroutineContext с помощью функции send или вне корутины с помощью функции trySend .

Internally, callbackFlow uses a channel , which is conceptually very similar to a blocking queue . A channel is configured with a capacity , the maximum number of elements that can be buffered. The channel created in callbackFlow has a default capacity of 64 elements. When you try to add a new element to a full channel, send suspends the producer until there's space for the new element, whereas trySend does not add the element to the channel and returns false immediately.

trySend немедленно добавляет указанный элемент в канал, только если это не нарушает его ограничения по вместимости, и затем возвращает успешный результат.

Дополнительные потоковые ресурсы