StateFlow ו-SharedFlow

‫StateFlow ו-SharedFlow הם Flow APIs שמאפשרים ל-Flows לשלוח עדכוני סטטוס בצורה אופטימלית ולשלוח ערכים לכמה צרכנים.

StateFlow

‫StateFlow הוא זרם נתונים ניתן לצפייה שמחזיק את המצב ופולט את עדכוני המצב הנוכחיים והחדשים למאספים שלו. אפשר לקרוא את ערך המצב הנוכחי גם דרך המאפיין value. כדי לעדכן את המצב ולשלוח אותו לתהליך, מקצים ערך חדש למאפיין value של המחלקה MutableStateFlow.

ב-Android, ‏ StateFlow מתאים מאוד למחלקות שצריכות לשמור על מצב ניתן לשינוי שניתן לצפייה.

בהמשך לדוגמאות מתוך Kotlin flows, אפשר לחשוף StateFlow מתוך LatestNewsViewModel כדי ש-View יוכל להאזין לעדכונים של מצב ממשק המשתמש, וכך מצב המסך ישמר באופן מובנה גם אחרי שינויים בהגדרות.

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

    // Backing property to avoid state updates from other classes
    private val _uiState = MutableStateFlow(LatestNewsUiState.Success(emptyList()))
    // The UI collects from this StateFlow to get its state updates
    val uiState: StateFlow<LatestNewsUiState> = _uiState

    init {
        viewModelScope.launch {
            newsRepository.favoriteLatestNews
                // Update UI with the latest favorite news
                // Writes to the value property of MutableStateFlow,
                // adding a new element to the flow and updating all
                // of its collectors
                .collect { favoriteNews ->
                    _uiState.value = LatestNewsUiState.Success(favoriteNews)
                }
        }
    }
}

// Represents different states for the LatestNews screen
sealed class LatestNewsUiState {
    data class Success(val news: List<ArticleHeadline>) : LatestNewsUiState()
    data class Error(val exception: Throwable) : LatestNewsUiState()
}

הכיתה שאחראית לעדכון MutableStateFlow היא המפיק, וכל הכיתות שאוספות מ-StateFlow הן הצרכנים. בניגוד לתהליך קר שנבנה באמצעות הכלי flow, תהליך StateFlow הוא פעיל: איסוף נתונים מהתהליך לא מפעיל קוד של יצרן. אובייקט StateFlow תמיד פעיל ונמצא בזיכרון, והוא הופך למועמד לאיסוף זבל רק כשאין הפניות אחרות אליו משורש איסוף הזבל.

כשצרכן חדש מתחיל לאסוף נתונים מהזרם, הוא מקבל את המצב האחרון בזרם ואת כל המצבים הבאים. אפשר לראות את ההתנהגות הזו במחלקות אחרות שאפשר לצפות בהן, כמו LiveData.

ה-View מאזין ל-StateFlow כמו בכל תהליך אחר:

class LatestNewsActivity : ComponentActivity() {
    private val latestNewsViewModel: LatestNewsViewModel = TODO() // getViewModel()

    override fun onCreate(savedInstanceState: Bundle?) {
        // ...
        // Start a coroutine in the lifecycle scope
        lifecycleScope.launch {
            // repeatOnLifecycle launches the block in a new coroutine every time the
            // lifecycle is in the STARTED state (or above) and cancels it when it's STOPPED.
            repeatOnLifecycle(Lifecycle.State.STARTED) {
                // Trigger the flow and start listening for values.
                // Note that this happens when lifecycle is STARTED and stops
                // collecting when the lifecycle is STOPPED
                latestNewsViewModel.uiState.collect { uiState ->
                    // New value received
                    when (uiState) {
                        is LatestNewsUiState.Success -> showFavoriteNews(uiState.news)
                        is LatestNewsUiState.Error -> showError(uiState.exception)
                    }
                }
            }
        }
    }
}

כדי להמיר כל זרימה ל-StateFlow, משתמשים באופרטור הביניים stateIn.

‫StateFlow,‏ Flow ו-LiveData

יש דמיון בין StateFlow לבין LiveData. שניהם הם מחלקות של בעלי נתונים שניתנים לצפייה, ושניהם פועלים לפי דפוס דומה כשמשתמשים בהם בארכיטקטורת האפליקציה.

עם זאת, חשוב לזכור שאופן הפעולה של StateFlow שונה מזה של LiveData:

  • ‫StateFlow דורש להעביר מצב ראשוני לקונסטרוקטור, אבל LiveData לא.
  • LiveData.observe() מבטל אוטומטית את הרישום של הצרכן כשהתצוגה עוברת למצב STOPPED, בעוד שאיסוף נתונים מ-StateFlow או מכל תהליך אחר לא מפסיק את האיסוף באופן אוטומטי. כדי להשיג את אותה התנהגות, צריך לאסוף את הרצף מLifecycle.repeatOnLifecycleבלוק.

הפיכת רצפי פעולות קרים לחמים באמצעות shareIn

‫StateFlow הוא זרם פעיל – הוא נשאר בזיכרון כל עוד הזרם נאסף או כל עוד קיימות הפניות אחרות אליו משורש מנגנון איסוף זבל. אפשר להפוך זרימות קרות לחמות באמצעות האופרטור shareIn.

לדוגמה, במקום שכל אוסף ייצור זרימה חדשה, אפשר לשתף את הנתונים שאוחזרו מ-Firestore בין האוספים באמצעות shareIn callbackFlow שנוצר בזרימות Kotlin. צריך להעביר את הפרטים הבאים:

  • CoroutineScope שמשמש לשיתוף של התהליך. ההיקף הזה צריך להיות פעיל למשך זמן ארוך יותר מכל צרכן, כדי שרצף הפעולות המשותף יישאר פעיל כל עוד צריך.
  • מספר הפריטים להפעלה חוזרת לכל מאסף חדש.
  • המדיניות בנושא התנהגות במהלך ההפעלה.

class NewsRemoteDataSource(
    private val externalScope: CoroutineScope
) {
    val latestNews: Flow<List<ArticleHeadline>> = flow {
        // ...
    }.shareIn(
        externalScope,
        replay = 1,
        started = SharingStarted.WhileSubscribed()
    )
}

בדוגמה הזו, הפריט האחרון שמועבר בזרם latestNews מועבר מחדש אל מאגר חדש, והזרם נשאר פעיל כל עוד externalScope פעיל ויש מאגרי נתונים פעילים. מדיניות ההתחלה של SharingStarted.WhileSubscribed() שומרת על פעילות של היצרן במעלה הזרם כל עוד יש מנויים פעילים. יש מדיניות הפעלה אחרת, כמו SharingStarted.Eagerly כדי להפעיל את המפיק באופן מיידי או SharingStarted.Lazily כדי להתחיל לשתף אחרי שהמנוי הראשון מופיע ולהשאיר את התהליך פעיל לנצח.

SharedFlow

הפונקציה shareIn מחזירה SharedFlow, זרם פעיל שפולט ערכים לכל הצרכנים שאוספים ממנו. ‫SharedFlow הוא הכללה של StateFlow שאפשר להגדיר אותה בצורה מפורטת.

אפשר ליצור SharedFlow בלי להשתמש ב-shareIn. לדוגמה, אפשר להשתמש ב-SharedFlow כדי לשלוח סימני זמן לשאר האפליקציה, כך שכל התוכן יתרענן באופן קבוע באותו הזמן. בנוסף לאחזור החדשות האחרונות, יכול להיות שתרצו לרענן את הקטע של פרטי המשתמש עם אוסף הנושאים המועדפים עליו. בקטע הקוד הבא, המחלקה TickHandler חושפת את המחלקה SharedFlow כדי שמחלקות אחרות יידעו מתי לרענן את התוכן שלה. בדומה ל-StateFlow, משתמשים במאפיין גיבוי מסוג MutableSharedFlow בכיתה כדי לשלוח פריטים לזרימה:

// Class that centralizes when the content of the app needs to be refreshed
class TickHandler(
    private val externalScope: CoroutineScope,
    private val tickIntervalMs: Long = 5000
) {
    // Backing property to avoid flow emissions from other classes
    private val _tickFlow = MutableSharedFlow<Unit>(replay = 0)
    val tickFlow: SharedFlow<Unit> = _tickFlow

    init {
        externalScope.launch {
            while (true) {
                _tickFlow.emit(Unit)
                delay(tickIntervalMs)
            }
        }
    }
}

class NewsRepository(
    // ...
    private val tickHandler: TickHandler,
    private val externalScope: CoroutineScope
) {
    init {
        externalScope.launch {
            // Listen for tick updates
            tickHandler.tickFlow.collect {
                refreshLatestNews()
            }
        }
    }

    suspend fun refreshLatestNews() { /* ... */ }
    // ...
}

אפשר להתאים אישית את ההתנהגות של SharedFlow בדרכים הבאות:

  • ‫replay מאפשרת לשלוח מחדש מספר ערכים שנשלחו בעבר למנויים חדשים.
  • באמצעות onBufferOverflow אפשר לציין מדיניות לגבי המצב שבו המאגר מלא בפריטים שצריך לשלוח. ערך ברירת המחדל הוא BufferOverflow.SUSPEND, שגורם למתקשר להשהות את השיחה. אפשרויות אחרות הן DROP_LATEST או DROP_OLDEST.

ל-MutableSharedFlow יש גם מאפיין subscriptionCount שמכיל את מספר האוספים הפעילים, כדי שתוכלו לבצע אופטימיזציה של הלוגיקה העסקית בהתאם. MutableSharedFlow כולל גם פונקציה resetReplayCache אם לא רוצים להפעיל מחדש את המידע האחרון שנשלח לזרימה.

מקורות מידע נוספים על תהליכים