Android RxJava Polling을 Kotlin Flow로 바꾸기: repeatOnLifecycle과 재시도
이 글은 기존 Android Retrofit2 RxJava polling: 일정 간격으로 API 호출하기와 짝을 이루는 마이그레이션 가이드입니다. RxJava로 구현된 polling을 Kotlin Flow와
repeatOnLifecycle로 옮기는 과정을 실행 가능한 코드로 다룹니다.
Android RxJava Polling을 Kotlin Flow로 바꾸기: repeatOnLifecycle과 재시도
이미 RxJava로 polling을 구현한 프로젝트에서 Kotlin Coroutines와 Flow로 옮기면 구조화된 동시성, 자연스러운 취소 처리, 생명 주기 연동을 더 단순하게 만들 수 있습니다. 이 글에서는 RxJava 3 polling 코드를 Kotlin Flow로 단계별로 바꾸고, 화면 생명 주기와 안전하게 연결하는 방법을 다룹니다.
| 핵심 질문 | 이 글의 답 |
|---|---|
RxJava interval은 Flow로 어떻게 바꾸나? |
flow + delay 또는 while().emit().delay() 루프 |
switchMapSingle은 Flow로 어떻게 표현하나? |
가장 최신 결과만 남기는 transformLatest / mapLatest |
| 오류로 polling이 끊기지 않게 하려면? | 요청 내부에서 예외를 상태로 변환, retryWhen으로 지수 백오프 |
| 화면이 사라지면 자동으로 멈추려면? | repeatOnLifecycle(STARTED)에서 수집 |
| 테스트는 어떻게? | runTest, TestDispatcher, 가짜 Repository |
왜 Flow로 바꾸는가
RxJava는 여전히 강력한 라이브러리지만, 새 Kotlin 프로젝트에서는 다음 이유로 Coroutines/Flow가 자연스러운 선택이 됩니다.
| 비교 | RxJava 3 | Kotlin Flow |
|---|---|---|
| 취소 전파 | Disposable 수동 관리 |
구조화된 동시성(부모-자식) |
| 생명 주기 연동 | LiveData/Lifecycle 래퍼 필요 |
repeatOnLifecycle 공식 API |
| 스레드 전환 | subscribeOn/observeOn |
Dispatchers와 FlowOn |
| 코틀린 친화도 | 자바 API 기반 | 언어 수준 지원 |
| 새 의존성 | RxJava, RxAndroid, Retrofit adapter | kotlinx-coroutines만 추가 |
이미 RxJava로 잘 동작하는 코드를 억지로 바꿀 필요는 없습니다. 다만 새 기능을 추가하거나 구조를 정리할 때 Flow로 옮기면 생명 주기 처리와 취소를 더 안전하게 만들 수 있습니다.
마이그레이션 전 RxJava 코드
출발점은 기존 글에서 다룬 RxJava polling ViewModel입니다. 핵심은 Observable.interval()과 switchMapSingle()입니다.
// Before: RxJava 3 polling
class PollingViewModel(
private val api: JokeApi
) : ViewModel() {
private val _uiState = MutableLiveData<PollingUiState>(PollingUiState.Idle)
val uiState: LiveData<PollingUiState> = _uiState
private var disposable: Disposable? = null
fun startPolling() {
if (disposable?.isDisposed == false) return
_uiState.value = PollingUiState.Loading
disposable = Observable.interval(0, 10, TimeUnit.SECONDS)
.switchMapSingle { tick ->
api.getRandomJoke()
.subscribeOn(Schedulers.io())
.map<PollingResult> { PollingResult.Success(it.value) }
.onErrorReturn { PollingResult.Failure(it.message ?: "오류") }
}
.observeOn(AndroidSchedulers.main())
.subscribe { result ->
_uiState.value = when (result) {
is PollingResult.Success -> PollingUiState.Success(result.message)
is PollingResult.Failure -> PollingUiState.Error(result.message)
}
}
}
fun stopPolling() {
disposable?.dispose()
disposable = null
}
override fun onCleared() {
stopPolling()
super.onCleared()
}
}
이 코드의 특징은 다음과 같습니다.
- 첫 요청을 즉시 실행하기 위해
interval(0, 10, ...)사용 - 이전 요청이 끝나지 않아도 새 요청이 들어오면 이전 요청을 취소하는
switchMapSingle - API 오류를 polling 종료가 아닌 상태 값으로 변환
Disposable을onCleared()에서 수동 정리
의존성 정리
Flow로 옮기려면 RxJava 관련 의존성을 줄이고 코루틴 의존성을 추가합니다. Retrofit은 suspend 함수를 지원하므로 Call adapter가 필요 없습니다.
dependencies {
// 코루틴과 수명 주기 연동
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-android:1.9.0")
implementation("androidx.lifecycle:lifecycle-runtime-ktx:2.11.0")
implementation("androidx.lifecycle:lifecycle-viewmodel-ktx:2.11.0")
// Retrofit은 suspend 함수 지원
implementation("com.squareup.retrofit2:retrofit:2.11.0")
implementation("com.squareup.retrofit2:converter-gson:2.11.0")
// RxJava는 점진적 제거 대상
// implementation("io.reactivex.rxjava3:rxjava:3.1.12")
// implementation("com.squareup.retrofit2:adapter-rxjava3:3.0.0")
}
정확한 버전은 글을 읽는 시점의 공식 문서에서 다시 확인하세요.
API 인터페이스 변경
RxJava Single 반환을 suspend 함수로 바꿉니다.
// Before
interface JokeApi {
@GET("jokes/random")
fun getRandomJoke(): Single<JokeResponse>
}
// After
interface JokeApi {
@GET("jokes/random")
suspend fun getRandomJoke(): JokeResponse
}
suspend 함수는 Retrofit이 백그라운드 디스패처로 자동 전환하므로 subscribeOn이 필요 없습니다.
Repository에서 polling Flow 만들기
API 호출을 감싸는 Repository를 만들고 polling Flow를 노출합니다. 핵심은 flow {} 빌더 안에서 루프를 돌며 결과를 방출하고 delay()로 주기를 만드는 것입니다.
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.isActive
import retrofit2.HttpException
import java.io.IOException
class JokeRepository(
private val api: JokeApi
) {
fun pollEvery(intervalMillis: Long = 10_000L): Flow<PollingUiState> = flow {
while (currentCoroutineContext().isActive) {
val state = try {
val response = api.getRandomJoke()
PollingUiState.Success(response.value)
} catch (error: IOException) {
PollingUiState.Error("네트워크 연결을 확인해 주세요.")
} catch (error: HttpException) {
PollingUiState.Error("서버 오류: ${error.code()}")
}
emit(state)
delay(intervalMillis)
}
}
}
| RxJava 연산자 | Flow 대응 |
|---|---|
Observable.interval() |
flow { while(active) { ... delay() } } |
switchMapSingle() |
transformLatest {} / mapLatest {} |
onErrorReturn() |
try/catch으로 상태 변환 |
subscribeOn(io()) |
suspend 함수 또는 flowOn(Dispatchers.IO) |
observeOn(mainThread()) |
수집 측에서 Dispatchers.Main |
주의: CancellationException을 삼키지 마세요
일반 catch (e: Exception)으로 모든 예외를 잡으면 코루틴 취소 신호까지 먹혀 화면이 종료되어도 polling이 멈추지 않을 수 있습니다. 네트워크 오류처럼 실제 처리할 예외만 구체적으로 잡으세요.
// 잘못된 예: 취소까지 무시
try {
emit(api.getRandomJoke())
} catch (e: Exception) {
emit(ErrorState)
}
// 올바른 예: 처리할 예외만 구체적으로
try {
emit(api.getRandomJoke())
} catch (e: IOException) {
emit(ErrorState)
}
최신 요청만 유지하기: transformLatest
RxJava switchMapSingle()은 새 값이 들어오면 이전 작업을 취소했습니다. Flow에서 같은 동작이 필요하면 transformLatest를 사용합니다.
import kotlinx.coroutines.flow.transformLatest
fun pollLatestOnly(
ticks: Flow<Long>,
intervalMillis: Long = 10_000L
): Flow<PollingUiState> = ticks.transformLatest { _ ->
val state = try {
PollingUiState.Success(api.getRandomJoke().value)
} catch (e: IOException) {
PollingUiState.Error("네트워크 오류")
}
emit(state)
kotlinx.coroutines.delay(intervalMillis)
}
| 상황 | 권장 연산자 |
|---|---|
| 최신 결과만 의미 | transformLatest / mapLatest |
| 모든 결과를 순서대로 처리 | transform 또는 일반 flow 루프 |
| 요청 중첩 허용 | 별도 제약 없이 루프에서 호출 |
대부분의 상태 polling은 "가장 최근 결과"만 보여주면 되므로 transformLatest가 자연스럽습니다.
재시도와 지수 백오프
요청 자체가 실패했을 때 polling을 종료하지 않고 잠시 기다렸다가 다시 시도하려면 retryWhen을 사용합니다.
import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.retryWhen
import kotlinx.coroutines.delay
fun pollWithBackoff(
maxAttempts: Int = 5,
baseDelayMillis: Long = 1_000L
): Flow<PollingUiState> = flow {
while (currentCoroutineContext().isActive) {
val response = api.getRandomJoke()
emit(PollingUiState.Success(response.value))
delay(10_000L)
}
}
.retryWhen { cause, attempt ->
if (attempt >= maxAttempts) {
emit(PollingUiState.Error("재시도 한도 초과: ${cause.message}"))
false
} else {
val backoff = baseDelayMillis * (1L shl attempt.toInt().coerceAtMost(6))
delay(backoff)
true
}
}
.catch { cause ->
emit(PollingUiState.Error("예상치 못한 오류: ${cause.message}"))
}
| 재시도 패턴 | 동작 |
|---|---|
| 고정 간격 | 매번 같은 시간 대기 |
| 지수 백오프 | 실패가 반복될수록 대기 시간 증가 |
| 최대 시도 횟수 | 초과 시 오류 상태 방출 후 종료 |
retryWhen은 업스트림 예외를 가로채서 재시도 여부를 결정합니다. 반면 내부에서 이미 상태로 변환한 예외는 retryWhen에 도달하지 않으므로 두 전략을 섞어 쓸 때는 흐름을 명확히 설계해야 합니다.
ViewModel에서 Flow 수집
ViewModel은 Repository Flow를 StateFlow로 변환해 UI에 노출합니다. viewModelScope를 사용하면 ViewModel이 삭제될 때 자동으로 취소됩니다.
import androidx.lifecycle.ViewModel
import androidx.lifecycle.viewModelScope
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.SharingStarted
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.stateIn
class PollingViewModel(
private val repository: JokeRepository
) : ViewModel() {
private val _uiState = MutableStateFlow<PollingUiState>(PollingUiState.Idle)
val uiState: StateFlow<PollingUiState> = _uiState.asStateFlow()
fun startPolling() {
viewModelScope.launch {
repository.pollEvery()
.collect { state -> _uiState.value = state }
}
}
}
Disposable 수동 정리가 사라졌습니다. viewModelScope가 취소되면 수집도 함께 취소됩니다.
| RxJava 패턴 | Flow 패턴 |
|---|---|
Disposable 필드 관리 |
Job 또는 스코프에 위임 |
onCleared()에서 dispose() |
viewModelScope 자동 취소 |
LiveData로 상태 노출 |
StateFlow로 상태 노출 |
Fragment에서 생명 주기 연결
화면이 보일 때만 수집하려면 repeatOnLifecycle을 사용합니다. 이 함수는 생명 주기 상태에 따라 수집을 시작하고 취소합니다.
import android.os.Bundle
import androidx.fragment.app.Fragment
import androidx.lifecycle.Lifecycle
import androidx.lifecycle.lifecycleScope
import androidx.lifecycle.repeatOnLifecycle
import kotlinx.coroutines.launch
class PollingFragment : Fragment() {
private val viewModel: PollingViewModel by viewModels()
override fun onViewCreated(view: View, savedInstanceState: Bundle?) {
super.onViewCreated(view, savedInstanceState)
viewLifecycleOwner.lifecycleScope.launch {
viewLifecycleOwner.repeatOnLifecycle(Lifecycle.State.STARTED) {
viewModel.uiState.collect { state ->
render(state)
}
}
}
viewModel.startPolling()
}
private fun render(state: PollingUiState) {
// 상태별 화면 갱신
}
}
| 생명 주기 상태 | 수집 동작 |
|---|---|
STARTED 이상 |
수집 시작 |
STARTED 미만 |
수집 취소, 새 값 대기하지 않음 |
| 화면 파괴 | viewLifecycleOwner 종료로 완전 정리 |
Lifecycle.State.CREATED에서 수집하면 onStop 이후에도 값을 처리하므로 의도하지 않은 갱신이 생길 수 있습니다. UI 갱신은 보통 STARTED 기준이 안전합니다.
Compose에서 수집
Jetpack Compose에서는 collectAsStateWithLifecycle()을 사용하면 생명 주기를 자동으로 처리합니다.
import androidx.compose.runtime.collectAsState
import androidx.lifecycle.compose.collectAsStateWithLifecycle
@Composable
fun PollingScreen(viewModel: PollingViewModel = viewModel()) {
val state by viewModel.uiState.collectAsStateWithLifecycle()
when (state) {
PollingUiState.Idle -> Text("대기 중")
PollingUiState.Loading -> CircularProgressIndicator()
is PollingUiState.Success -> Text(state.message)
is PollingUiState.Error -> Text(state.message)
}
}
collectAsStateWithLifecycle()은 lifecycle-runtime-compose 아티팩트가 필요합니다.
implementation("androidx.lifecycle:lifecycle-runtime-compose:2.11.0")
테스트
Flow 테스트는 runTest와 TestDispatcher로 시간을 제어합니다. 가짜 Repository로 폴링 흐름을 검증합니다.
import app.cash.turbine.test
import kotlinx.coroutines.test.runTest
import kotlin.test.assertEquals
class JokeRepositoryTest {
@Test
fun 첫_값은_즉시_방출된다() = runTest {
val repository = JokeRepository(FakeJokeApi(listOf("a", "b")))
repository.pollEvery(intervalMillis = 1_000L).test {
assertEquals(PollingUiState.Success("a"), awaitItem())
assertEquals(PollingUiState.Success("b"), awaitItem())
awaitComplete() // 가짜 API가 비면 예외 또는 종료 시나리오에 따라
}
}
}
private class FakeJokeApi(private val responses: List<String>) : JokeApi {
private val queue = ArrayDeque(responses)
override suspend fun getRandomJoke(): JokeResponse {
return JokeResponse(queue.removeFirst())
}
}
Turbine(app.cash.turbine:turbine)은 Flow 테스트를 위한 인기 라이브러리입니다. 시간 기반 동작은 runTest의 가상 시간 제어와 advanceTimeBy로 검증합니다.
| 테스트 도구 | 용도 |
|---|---|
runTest |
코루틴 시간 가속 |
TestDispatcher |
디스패처 제어 |
| Turbine | Flow 항목 순차 검증 |
| 가짜 Repository | API 없이 흐름 검증 |
마이그레이션 체크리스트
RxJava에서 Flow로 옮길 때 점검할 항목입니다.
- [ ] API 인터페이스가
suspend함수로 변경됨 - [ ] RxJava 의존성이 더 이상 사용되지 않으면 제거
- [ ] polling Flow가 오류를 상태로 변환함
- [ ]
CancellationException을 일반 예외로 잡지 않음 - [ ] 화면 수집이
repeatOnLifecycle또는collectAsStateWithLifecycle로 연결됨 - [ ] ViewModel이
viewModelScope로 자동 취소됨 - [ ] 단위 테스트가
runTest+ Turbine으로 작성됨 - [ ] 지수 백오프 또는 최대 재시도가 필요한 경우
retryWhen적용됨
마무리
RxJava polling을 Kotlin Flow로 옮기면 취소 처리와 생명 주기 연동을 언어 수준에서 더 안전하게 다룰 수 있습니다. flow {} 루프로 interval을, transformLatest로 switchMapSingle을, retryWhen으로 재시도를 대체하면 대부분의 polling 시나리오를 자연스럽게 표현할 수 있습니다.
중요한 것은 화면이 사라졌을 때 polling이 멈추도록 repeatOnLifecycle이나 collectAsStateWithLifecycle을 일관되게 사용하는 것입니다. 구조화된 동시성이 이 보장을 프레임워크 수준에서 도와줍니다.
댓글
댓글 쓰기