Android'de Kotlin akışı

Coroutine'larda flow, yalnızca tek bir değer döndüren askıya alma işlevlerinin aksine, birden fazla değeri sırayla yayınlayabilen bir türdür. Örneğin, bir akışı kullanarak veritabanından canlı güncellemeler alabilirsiniz.

Akışlar, eş yordamların üzerine kurulur ve birden fazla değer sağlayabilir. Akış, kavramsal olarak eşzamansız şekilde hesaplanabilen bir veri akışıdır. Yayılan değerler aynı türde olmalıdır. Örneğin, Flow<Int>, tam sayı değerleri veren bir akıştır.

Akış, bir değer dizisi oluşturan Iterator'a çok benzer ancak değerleri eşzamansız olarak üretmek ve kullanmak için askıya alma işlevlerini kullanır. Bu, örneğin akışın ana iş parçacığını engellemeden bir sonraki değeri üretmek için güvenli bir şekilde ağ isteğinde bulunabileceği anlamına gelir.

Veri akışlarında üç öğe yer alır:

  • Üretici, akışa eklenen veriler üretir. Coroutines sayesinde akışlar, verileri eşzamansız olarak da üretebilir.
  • (İsteğe bağlı) Aracı kuruluşlar, akışa gönderilen her değeri veya akışın kendisini değiştirebilir.
  • Bir tüketici, akıştaki değerleri tüketir.

Veri akışlarına dahil olan öğeler; tüketici, isteğe bağlı aracılar ve üretici
Şekil 1. Veri akışlarına dahil olan öğeler: tüketici, isteğe bağlı aracı ve üretici.

Android'de depo genellikle yerel veritabanlarından (Room) veya bulut arka uçlarından (ör. Firebase) kullanıcı arayüzü verilerini yayınlar. Bu durumda kullanıcı arayüzü, verilerin tüketicisi olarak hareket eder. Aksine, kullanıcı arayüzü diğer katmanların kullanması için kullanıcı girişi etkinlikleri oluşturabilir. Ara katmanlar (ör. ViewModel örnekleri) akışı her katmanın gereksinimlerini karşılayacak şekilde dönüştürür.

Akış oluşturma

Akış oluşturmak için akış oluşturma aracı API'lerini kullanın. flow oluşturucu işlevi, emit işlevini kullanarak veri akışına manuel olarak yeni değerler gönderebileceğiniz yeni bir akış oluşturur.

Aşağıdaki örnekte, bir veri kaynağı en son haberleri sabit aralıklarla otomatik olarak getirir. Askıya alma işlevi, birden fazla ardışık değer döndüremediğinden veri kaynağı, bu koşulu karşılamak için bir akış oluşturup döndürür. Bu durumda veri kaynağı, üretici olarak hareket eder.

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 oluşturucu, bir eş yordam içinde yürütülür. Bu nedenle, aynı asenkron API'lerden yararlanır ancak bazı kısıtlamalar geçerlidir:

  • Akışlar sıralıdır. Üretici bir eş yordamda olduğundan, bir askıya alma işlevi çağrıldığında üretici, askıya alma işlevi döndürülene kadar askıya alınır. Örnekte, yapımcı fetchLatestNews ağ isteği tamamlanana kadar yürütmeyi duraklatır. Sonuç ancak bu durumda akışa gönderilir.
  • flow oluşturucu ile üretici, farklı bir CoroutineContext değerini emit yapamaz. Bu nedenle, yeni eş yordamlar oluşturarak veya withContext kod bloklarını kullanarak emit işlevini farklı bir CoroutineContext içinde çağırmayın. Bu durumlarda, callbackFlow gibi diğer akış oluşturucuları kullanabilirsiniz.

Yayını değiştirme

Aracılar, değerleri kullanmadan veri akışını değiştirmek için aracı operatörleri kullanabilir. Bu operatörler, bir veri akışına uygulandığında değerler gelecekte kullanılana kadar yürütülmeyen bir işlem zinciri oluşturan işlevlerdir. Ara operatörler hakkında daha fazla bilgiyi Akış referans belgelerinde bulabilirsiniz.

Aşağıdaki örnekte, depo katmanı, View üzerinde gösterilecek verileri dönüştürmek için ara operatör map kullanır:

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

Ara operatörler, akışa bir öğe gönderildiğinde tembelce yürütülen bir işlem zinciri oluşturarak art arda uygulanabilir. Bir akışa yalnızca ara operatör uygulamanın akış toplama işlemini başlatmadığını unutmayın.

Bir akıştan toplama

Değerleri dinlemeye başlamak için akışı tetiklemek üzere terminal operatörü kullanın. Akışta yayınlanan tüm değerleri almak için collect kullanın. Terminal operatörleri hakkında daha fazla bilgiyi resmi akış belgelerinde bulabilirsiniz.

collect, askıya alma işlevi olduğundan bir coroutine içinde yürütülmesi gerekir. Her yeni değerde çağrılan bir lambda'yı parametre olarak alır. Askıya alma işlevi olduğundan, collect işlevini çağıran eş yordam, akış kapatılana kadar askıya alınabilir.

Önceki örnekten devam ederek, depo katmanındaki verileri kullanan bir ViewModel öğesinin basit bir uygulamasını aşağıda bulabilirsiniz:

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

Akışın toplanması, en son haberleri yenileyen ve ağ isteğinin sonucunu sabit bir aralıkta yayınlayan üreticiyi tetikler. Üretici, while(true) döngüsüyle her zaman etkin kaldığından ViewModel temizlendiğinde ve viewModelScope iptal edildiğinde veri akışı kapatılır.

Akış toplama işlemi aşağıdaki nedenlerle durabilir:

  • Önceki örnekte gösterildiği gibi, toplayan coroutine iptal edilir. Bu işlem, temel üreticiyi de durdurur.
  • Üretici, öğe yayınlamayı tamamlar. Bu durumda veri akışı kapatılır ve collect işlevini çağıran eş yordamın yürütülmesine devam edilir.

Diğer ara operatörlerle belirtilmediği sürece akışlar soğuk ve tembeldir. Bu, akışta her terminal operatörü çağrıldığında üretici kodunun yürütüldüğü anlamına gelir. Önceki örnekte, birden fazla akış toplayıcının olması, veri kaynağının en son haberleri farklı sabit aralıklarla birden çok kez getirmesine neden olur. Birden fazla tüketici aynı anda veri topladığında bir akışı optimize etmek ve paylaşmak için shareIn operatörünü kullanın.

Beklenmeyen istisnaları yakalama

Üreticinin uygulanması, üçüncü taraf kitaplığından gelebilir. Bu, beklenmedik istisnalar oluşturabileceği anlamına gelir. Bu istisnaları işlemek için catch ara operatörünü kullanın.

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

Önceki örnekte, bir istisna oluştuğunda yeni bir öğe alınmadığı için collect lambda çağrılmaz.

catch, akışa emit öğeleri de ekleyebilir. Örnek depolama alanı katmanı, bunun yerine önbelleğe alınmış değerleri emit olabilir:

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

Bu örnekte, bir istisna oluştuğunda akışa istisna nedeniyle yeni bir öğe gönderildiğinden collect lambda'sı çağrılır.

Farklı bir CoroutineContext'te yürütme

Varsayılan olarak, bir flow oluşturucunun üreticisi, kendisinden toplayan eş yordamın CoroutineContext içinde yürütülür ve daha önce belirtildiği gibi, farklı bir CoroutineContext değerleri emit olamaz. Bu davranış bazı durumlarda istenmeyebilir. Örneğin, bu konu boyunca kullanılan örneklerde, depo katmanı viewModelScope tarafından kullanılan Dispatchers.Main üzerinde işlemler yapmamalıdır.

Bir akışın CoroutineContext değerini değiştirmek için ara operatörü flowOn kullanın. flowOn, yukarı akışın CoroutineContext değerini değiştirir. Bu, üreticinin ve flowOn öncesinde (veya üstünde) uygulanan tüm ara operatörlerin değiştiği anlamına gelir. Aşağı akış (tüketiciyle birlikte flowOn sonraki ara operatörler) etkilenmez ve akıştan collect için kullanılan CoroutineContext üzerinde yürütülür. Birden fazla flowOn operatörü varsa her biri akışı mevcut konumundan değiştirir.

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

Bu kodla, onEach ve map operatörleri defaultDispatcher kullanırken catch operatörü ve tüketici, viewModelScope tarafından kullanılan Dispatchers.Main üzerinde yürütülür.

Veri kaynağı katmanı G/Ç çalışması yaptığından G/Ç işlemleri için optimize edilmiş bir gönderici kullanmanız gerekir:

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

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

Jetpack kitaplıklarındaki ve veri kaynaklarındaki akışlar

Flow, birçok Jetpack kitaplığına entegre edilmiştir ve Android üçüncü taraf kitaplıkları ile bulut SDK'larında popülerdir. Flow, canlı veri güncellemeleri ve sonsuz veri akışları için idealdir.

Veritabanındaki değişikliklerden haberdar olmak için Oda ile Akış'ı kullanabilirsiniz. Veri erişimi nesnelerini (DAO) kullanırken anlık güncellemeleri almak için Flow türünü döndürün.

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

Example tablosunda her değişiklik olduğunda, veritabanındaki yeni öğeleri içeren yeni bir liste yayınlanır.

Benzer şekilde, Cloud Firestore gibi anlık bulut veri kaynakları, canlı veri güncellemelerini doğrudan depo katmanınıza aktaran Flow uzantılar sağlar:

class RemoteMessagesDataSource(
    private val firestore: FirebaseFirestore
) {
    // Returns a Flow emitting updated message lists when collection changes
    val messages: Flow<List<Message>> =
        firestore.collection("messages").dataObjects()
}

Bu, yerel depolama ile doğal olarak eşleşir: Deponuz, hem Room'daki (yerel önbellek) hem de Firebase'deki (uzak bulut verileri) akışları birleştirerek çevrimdışı öncelikli ve anlık bir kullanıcı deneyimi sağlayabilir.

Geri çağırmaya dayalı API'leri akışlara dönüştürme

callbackFlow, geri çağırmaya dayalı API'leri akışlara dönüştürmenize olanak tanıyan bir akış oluşturucudur. Bu, özellikle Cloud Firestore veya Firebase Remote Config gibi SDK'larda birden fazla etkinlik dinleyicisiyle çalışırken kullanışlıdır.

Geri çağırmaya dayalı bir dinleyici API'yi akışa dönüştürmek için (ör. Firestore veritabanı güncellemelerini dinlemek) aşağıdaki kodu kullanın:

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 oluşturucunun aksine, callbackFlow, send işleviyle farklı bir CoroutineContext'den veya trySend işleviyle bir eşzamanlılık rutininin dışından değerlerin yayınlanmasına olanak tanır.

callbackFlow, dahili olarak bir kanal kullanır. Bu kanal, kavramsal olarak engelleme kuyruğuna çok benzer. Bir kanal, arabelleğe alınabilecek maksimum öğe sayısı olan kapasite ile yapılandırılır. callbackFlow içinde oluşturulan kanalın varsayılan kapasitesi 64 öğedir. Dolu bir kanala yeni bir öğe eklemeye çalıştığınızda send, yeni öğe için yer açılana kadar üreticiyi askıya alır. trySend ise öğeyi kanala eklemez ve hemen false döndürür.

trySend, belirtilen öğeyi kapasite kısıtlamalarını ihlal etmemesi koşuluyla kanala hemen ekler ve başarılı sonucu döndürür.

Ek akış kaynakları