You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何优化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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 19:53:16