如何优化Kotlin Flow事件处理,适配网络状态变更场景?
问题描述
现有一个处理事件的Manager类,当网络状态变更但事件仍在处理时表现异常。原代码如下:
初始化代码
fun init() { job = CoroutineScope(ioDispatcher).launch { combine( networkMonitor.isOnline, fcmEventsRepository.getPendingFcmEvent() ) { isOnline, event -> isOnline to event }.collect { (isOnline, event) -> if (isOnline && event != null) { processEvent(event) } } } }
事件处理函数
private suspend fun processEvent(event: FcmEvent) = withContext(ioDispatcher) { var success = false var retries = 0 while (!success && isActive) { try { val isProcessed = fcmEventsRepository.isEventProcessed(event.id) if (isProcessed) return@withContext retryWithExponentialBackoff( maxRetries = 5, initialDelayMillis = 1000, factor = 2, action = { syncRepository.syncByEventType(event.type) fcmEventsRepository.updateFcmEventStatus( event.id, ProcessingStatus.Processed ) } ) success = true } catch (e: Exception) { Timber.e(e, "Failed to process event $event, retrying in 1 minute...") delay(eventCheckDelay) retries++ } } }
需要满足以下要求:
- 同一时间仅处理单个事件
- 事件处理过程中设备失去连接时,取消处理
- 网络状态变更时,同一事件不会被重复处理
解决方案
1. 确保同一时间仅处理单个事件
给Manager类加一个互斥锁,每次调用事件处理函数时用锁包裹,强制同一时间只能执行一个事件处理流程:
// 在Manager类中声明互斥锁 private val processingMutex = Mutex() // 修改init中的collect逻辑 collect { (isOnline, event) -> if (isOnline && event != null) { processingMutex.withLock { processEvent(event) } } }
2. 网络断开时取消事件处理
把原有的combine换成flatMapLatest,网络离线时直接发射空值,触发collectLatest取消当前正在处理的协程;同时在processEvent的循环中,每次都检查网络状态,确保离线后立即停止重试:
// 修改init函数 fun init() { job = CoroutineScope(ioDispatcher).launch { networkMonitor.isOnline.flatMapLatest { isOnline -> if (isOnline) { fcmEventsRepository.getPendingFcmEvent() } else { // 网络离线时发射空值,取消当前处理流程 flowOf(null) } }.collectLatest { event -> event?.let { processingMutex.withLock { processEvent(it) } } } } } // 修改processEvent的循环条件 private suspend fun processEvent(event: FcmEvent) = withContext(ioDispatcher) { var success = false var retries = 0 // 循环中同时检查协程活跃状态和网络状态 while (!success && isActive && networkMonitor.isOnline.first()) { try { // ...原有处理逻辑 } catch (e: Exception) { Timber.e(e, "Failed to process event $event, retrying in 1 minute...") delay(eventCheckDelay) retries++ // 重试前再次确认网络状态 if (!networkMonitor.isOnline.first()) break } } }
3. 避免网络状态变更时重复处理同一事件
给事件新增Processing中间状态,开始处理前就把事件标记为该状态,同时让仓库的getPendingFcmEvent只返回Pending状态的事件,彻底避免重复获取正在处理的事件;处理失败时再把状态改回Pending,方便后续重试:
// 修改processEvent函数 private suspend fun processEvent(event: FcmEvent) = withContext(ioDispatcher) { var success = false var retries = 0 // 先标记为处理中,防止被重复获取 fcmEventsRepository.updateFcmEventStatus(event.id, ProcessingStatus.Processing) try { while (!success && isActive && networkMonitor.isOnline.first()) { // 双重检查事件状态,避免其他流程已处理该事件 val currentStatus = fcmEventsRepository.getEventStatus(event.id) if (currentStatus == ProcessingStatus.Processed) { success = true break } retryWithExponentialBackoff( maxRetries = 5, initialDelayMillis = 1000, factor = 2, action = { syncRepository.syncByEventType(event.type) fcmEventsRepository.updateFcmEventStatus( event.id, ProcessingStatus.Processed ) } ) success = true } } catch (e: Exception) { Timber.e(e, "Failed to process event $event, retrying in 1 minute...") delay(eventCheckDelay) retries++ } finally { // 未成功处理时,重置为待处理状态 if (!success) { fcmEventsRepository.updateFcmEventStatus(event.id, ProcessingStatus.Pending) } } } // 仓库层补充逻辑:只返回Pending状态的事件 fun getPendingFcmEvent(): Flow<FcmEvent?> { return db.fcmEventsDao().getEventsByStatus(ProcessingStatus.Pending) .map { it.firstOrNull()?.toDomainModel() } }
内容的提问来源于stack exchange,提问作者Emmanuel Mtali
相关产品推荐
相关产品推荐

