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 bulunur:
- Ü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ıştan gelen değerleri tüketir.
Android'de depo, genellikle kullanıcı arayüzünü (UI) tüketici olarak kullanan ve verileri nihai olarak görüntüleyen bir kullanıcı arayüzü veri üreticisidir. Diğer zamanlarda ise kullanıcı arayüzü katmanı, kullanıcı girişi etkinliklerinin üreticisi olur ve hiyerarşinin diğer katmanları bunları tüketir. Üretici ile tüketici arasındaki katmanlar genellikle, veri akışını bir sonraki katmanın gereksinimlerine uyacak şekilde değiştiren aracı olarak hareket eder.
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, 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, 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ı
fetchLatestNewsağ isteği tamamlanana kadar askıya alınır. Sonuç ancak bu durumda akışa gönderilir. flowoluşturucu ile üretici, farklı birCoroutineContextdeğeriniemityapamaz. Bu nedenle, yeni coroutine'ler oluşturarak veyawithContextkod bloklarını kullanarakemitişlevini farklı birCoroutineContextiçinde çağırmayın. Bu durumlarda,callbackFlowgibi diğer akış oluşturucuları kullanabilirsiniz.
Yayını değiştirme
Aracılar, değerleri kullanmadan veri akışını değiştirmek için ara 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ıyor:
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şturacak şekilde art arda uygulanabilir. Bir akışa yalnızca ara operatör uygulamanın akış toplama işlemini başlatmadığını unutmayın.
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 bir 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ığı için 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
collectiş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 bir terminal operatörü her ç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 fazla 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 nedenle, beklenmedik istisnalar oluşturabilir. 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'sı ç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'dan emit değerleri alamaz. 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, yukarı 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 akışlar
Flow, birçok Jetpack kitaplığına entegre edilmiştir ve Android üçüncü taraf kitaplıkları arasında popülerdir. Flow, canlı veri güncellemeleri ve sonsuz veri akışları için idealdir.
Bir veritabanındaki değişikliklerden haberdar olmak için Oda ile Akış'ı kullanabilirsiniz. Veri erişimi nesnelerini (DAO) kullanırken anlık güncellemeler 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.
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.
Örneğin, Firebase Firestore Android API'leri geri çağırmaları kullanır.
Bu API'leri akışlara dönüştürmek ve Firestore veritabanı güncellemelerini dinlemek için aşağıdaki kodu kullanabilirsiniz:
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. Tam 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.