Kotlin na Androidzie

W korutynach przepływ to typ, który może emitować wiele wartości sekwencyjnie, w przeciwieństwie do funkcji zawieszających, które zwracają tylko jedną wartość. Możesz na przykład użyć przepływu, aby otrzymywać na żywo aktualizacje z bazy danych.

Przepływy są oparte na korutynach i mogą dostarczać wiele wartości. Przepływ to koncepcyjnie strumień danych, który można obliczać asynchronicznie. Emitowane wartości muszą być tego samego typu. Na przykład Flow<Int> to przepływ, który emituje wartości całkowite.

Przepływ jest bardzo podobny do Iterator, który generuje sekwencję wartości, ale używa funkcji zawieszających do generowania i używania wartości asynchronicznie. Oznacza to na przykład, że przepływ może bezpiecznie wysłać żądanie sieciowe, aby wygenerować następną wartość, bez blokowania głównego wątku.

W strumieniach danych występują 3 rodzaje elementów:

  • Producent generuje dane, które są dodawane do strumienia. Dzięki korutynom automatyzacje mogą też generować dane asynchronicznie.
  • (Opcjonalnie) Pośrednicy mogą modyfikować każdą wartość emitowaną do strumienia lub sam strumień.
  • Konsument wykorzystuje wartości ze strumienia.

podmioty uczestniczące w strumieniach danych: konsument, pośrednicy (opcjonalnie) i producent;
Rysunek 1. Podmioty zaangażowane w strumienie danych: konsument, opcjonalni pośrednicy i producent.

W Androidzie repozytorium jest zwykle producentem danych interfejsu, którego odbiorcą jest interfejs użytkownika, który ostatecznie wyświetla dane. Czasami warstwa interfejsu jest producentem zdarzeń wejściowych użytkownika, a inne warstwy hierarchii je wykorzystują. Warstwy między producentem a konsumentem zwykle pełnią rolę pośredników, którzy modyfikują strumień danych, aby dostosować go do wymagań kolejnej warstwy.

Tworzenie przepływu

Aby tworzyć automatyzacje, użyj interfejsów API kreatora automatyzacji. Funkcja flow kreatora tworzy nowy przepływ, w którym możesz ręcznie emitować nowe wartości do strumienia danych za pomocą funkcji emit.

W przykładzie poniżej źródło danych automatycznie pobiera najnowsze wiadomości w ustalonych odstępach czasu. Funkcja zawieszająca nie może zwracać wielu kolejnych wartości, więc źródło danych tworzy i zwraca przepływ, aby spełnić to wymaganie. W tym przypadku źródło danych pełni rolę producenta.

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>
}

Funkcja flow jest wykonywana w ramach korutyny. Dlatego korzysta z tych samych asynchronicznych interfejsów API, ale obowiązują pewne ograniczenia:

  • Automatyzacje są sekwencyjne. Ponieważ producent znajduje się w współprogramie, podczas wywoływania funkcji zawieszającej zawiesza się do momentu, aż funkcja zawieszająca zwróci wartość. W tym przykładzie producent wstrzymuje działanie do momentu zakończenia żądania sieciowego fetchLatestNews. Dopiero wtedy wynik jest przesyłany do strumienia.
  • W przypadku narzędzia do tworzenia flow producent nie może emit wartości z innego CoroutineContext. Dlatego nie wywołuj funkcji emit w innejCoroutineContext, tworząc nowe współprogramy lub używając withContextbloków kodu. W takich przypadkach możesz użyć innych narzędzi do tworzenia przepływów, np. callbackFlow.

Modyfikowanie strumienia

Pośrednicy mogą używać operatorów pośrednich do modyfikowania strumienia danych bez wykorzystywania wartości. Te operatory to funkcje, które po zastosowaniu do strumienia danych tworzą łańcuch operacji, które nie są wykonywane, dopóki wartości nie zostaną wykorzystane w przyszłości. Więcej informacji o operatorach pośrednich znajdziesz w dokumentacji referencyjnej dotyczącej przepływów.

W przykładzie poniżej warstwa repozytorium używa operatora pośredniego map do przekształcania danych, które mają być wyświetlane w 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) }
}

Operatory pośrednie można stosować jeden po drugim, tworząc łańcuch operacji, które są wykonywane w sposób odroczony, gdy element jest emitowany do przepływu. Pamiętaj, że samo zastosowanie operatora pośredniego do strumienia nie rozpoczyna zbierania przepływu.

Zbieranie danych z przepływu

Użyj operatora terminala, aby wywołać przepływ i rozpocząć nasłuchiwanie wartości. Aby uzyskać wszystkie wartości w strumieniu w momencie ich wyemitowania, użyj collect. Więcej informacji o operatorach terminali znajdziesz w oficjalnej dokumentacji przepływu.

Ponieważ collect jest funkcją zawieszającą, musi być wykonywana w współprogramie. Przyjmuje jako parametr funkcję lambda, która jest wywoływana dla każdej nowej wartości. Ponieważ jest to funkcja zawieszająca, współprogram, który wywołuje collect, może zostać zawieszony do momentu zamknięcia przepływu.

Kontynuując poprzedni przykład, oto prosta implementacja ViewModel, która korzysta z danych z warstwy repozytorium:

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
            }
        }
    }
}

Zbieranie danych wywołuje producenta, który odświeża najnowsze wiadomości i emituje wynik żądania sieciowego w stałych odstępach czasu. Producent pozostaje zawsze aktywny w pętli while(true), więc strumień danych zostanie zamknięty, gdy ViewModel zostanie wyczyszczony, a viewModelScope zostanie anulowany.

Zbieranie danych o przepływach może zostać zatrzymane z tych powodów:

  • Współprogram, który zbiera dane, jest anulowany, jak pokazano w poprzednim przykładzie. Spowoduje to również zatrzymanie bazowego producenta.
  • Producent kończy emitowanie elementów. W takim przypadku strumień danych jest zamykany, a współprogram, który wywołał funkcję collect, wznawia wykonanie.

Strumienie są zimneleniwe, chyba że określono inaczej za pomocą innych operatorów pośrednich. Oznacza to, że kod producenta jest wykonywany za każdym razem, gdy w przepływie wywoływany jest operator końcowy. W poprzednim przykładzie posiadanie wielu kolektorów przepływu powoduje, że źródło danych pobiera najnowsze wiadomości wielokrotnie w różnych stałych odstępach czasu. Aby zoptymalizować i udostępnić przepływ, gdy wielu konsumentów zbiera dane w tym samym czasie, użyj operatora shareIn.

Wyłapywanie nieoczekiwanych wyjątków

Implementacja producenta może pochodzić z biblioteki innej firmy. Oznacza to, że może zgłaszać nieoczekiwane wyjątki. Aby obsłużyć te wyjątki, użyj operatora pośredniego 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
                }
        }
    }
}

W poprzednim przykładzie, gdy wystąpi wyjątek, funkcja lambda collect nie zostanie wywołana, ponieważ nie otrzymano nowego elementu.

catch może też emit elementy do przepływu. Przykładowa warstwa repozytorium mogłaby emit wartości z pamięci podręcznej:

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()) }
}

W tym przykładzie, gdy wystąpi wyjątek, wywoływana jest funkcja lambda collect, ponieważ w strumieniu pojawił się nowy element z powodu wyjątku.

Wykonywanie w innym elemencie CoroutineContext

Domyślnie producent flow wykonuje się w CoroutineContext współprogramu, który z niego zbiera, i jak wspomnieliśmy wcześniej, nie może emit wartości z innego CoroutineContext. W niektórych przypadkach takie działanie może być niepożądane. Na przykład w przykładach użytych w tym temacie warstwa repozytorium nie powinna wykonywać operacji na Dispatchers.Main, które są używane przez viewModelScope.

Aby zmienić CoroutineContext przepływu, użyj operatora pośredniego flowOn. flowOn zmienia CoroutineContext przepływu nadrzędnego, co oznacza, że producent i wszyscy operatorzy pośredni są stosowani przed (lub powyżej) flowOn. Proces downstream (operatorzy pośredni po flowOn wraz z konsumentem) nie jest objęty tym działaniem i jest wykonywany na CoroutineContext używanym do collect z procesu. Jeśli jest kilka operatorów flowOn, każdy z nich zmienia element nadrzędny z jego bieżącej lokalizacji.

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())
            }
}

W tym kodzie operatory onEachmap używają defaultDispatcher, a operator catch i konsument są wykonywane na Dispatchers.Main używanym przez viewModelScope.

Ponieważ warstwa źródła danych wykonuje operacje wejścia-wyjścia, należy użyć dyspozytora zoptymalizowanego pod kątem takich operacji:

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

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

Flow w bibliotekach Jetpack

Flow jest zintegrowany z wieloma bibliotekami Jetpack i jest popularny wśród bibliotek innych firm na Androida. Flow doskonale sprawdza się w przypadku aktualizacji danych na żywo i niekończących się strumieni danych.

Możesz używać Flow z Room, aby otrzymywać powiadomienia o zmianach w bazie danych. Jeśli używasz obiektów umożliwiających dostęp do danych (DAO), zwracaj typ Flow, aby otrzymywać aktualizacje na żywo.

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

Za każdym razem, gdy w tabeli Example nastąpi zmiana, emitowana jest nowa lista z nowymi elementami w bazie danych.

Przekształcanie interfejsów API opartych na wywołaniach zwrotnych w funkcje typu flow

callbackFlow to narzędzie do tworzenia przepływów, które umożliwia przekształcanie interfejsów API opartych na wywołaniach zwrotnych w przepływy. Na przykład interfejsy API Firebase Firestore na Androida używają wywołań zwrotnych.

Aby przekształcić te interfejsy API w przepływy i nasłuchiwać aktualizacji bazy danych Firestore, możesz użyć tego kodu:

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() }
    }
}

W przeciwieństwie do konstruktora flow konstruktor callbackFlow umożliwia emitowanie wartości z innego CoroutineContext za pomocą funkcji send lub poza korutyną za pomocą funkcji trySend.

Wewnętrznie callbackFlow używa kanału, który jest koncepcyjnie bardzo podobny do blokującej kolejki. Kanał jest skonfigurowany z pojemnością, czyli maksymalną liczbą elementów, które można buforować. Kanał utworzony w callbackFlow ma domyślną pojemność 64 elementów. Gdy spróbujesz dodać nowy element do pełnego kanału, send wstrzyma producenta, dopóki nie będzie miejsca na nowy element, a trySend nie doda elementu do kanału i natychmiast zwróci false.

trySend natychmiast dodaje określony element do kanału, ale tylko wtedy, gdy nie narusza to ograniczeń jego pojemności, a następnie zwraca wynik potwierdzający.

Dodatkowe zasoby dotyczące przepływu