Android RxJava Polling을 Kotlin Flow로 바꾸기: repeatOnLifecycle과 재시도

이 글은 기존 Android Retrofit2 RxJava polling: 일정 간격으로 API 호출하기와 짝을 이루는 마이그레이션 가이드입니다. RxJava로 구현된 polling을 Kotlin Flow와 repeatOnLifecycle로 옮기는 과정을 실행 가능한 코드로 다룹니다.

RxJava polling 연산자를 Kotlin Flow의 책임으로 옮기는 대응 관계

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 DispatchersFlowOn
코틀린 친화도 자바 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 종료가 아닌 상태 값으로 변환
  • DisposableonCleared()에서 수동 정리

의존성 정리

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을 사용합니다. 이 함수는 생명 주기 상태에 따라 수집을 시작하고 취소합니다.

Fragment 생명 주기에 따라 polling Flow 수집이 시작되고 취소되는 흐름

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 테스트는 runTestTestDispatcher로 시간을 제어합니다. 가짜 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을, transformLatestswitchMapSingle을, retryWhen으로 재시도를 대체하면 대부분의 polling 시나리오를 자연스럽게 표현할 수 있습니다.

중요한 것은 화면이 사라졌을 때 polling이 멈추도록 repeatOnLifecycle이나 collectAsStateWithLifecycle을 일관되게 사용하는 것입니다. 구조화된 동시성이 이 보장을 프레임워크 수준에서 도와줍니다.


참고 자료

댓글

이 블로그의 인기 게시물

pyautogui 예제 모니터 특정 위치 색상 구하고 비교해서 클릭 이벤트 하기

vscode 에서 WSL 개발환경 동작하지 않는 경우

React에서 Socket.IO Client 연결하기: CORS와 useEffect 정리