Kotlin Flow消费者未接收全部事件问题求助
Kotlin生产者/消费者模式异常排查:消费者接收事件数量不符
问题背景
我在实现Kotlin生产者/消费者模式时遇到问题:方法接收需发送至AWS服务的元素列表,设计了三个组件:
- 消费者:读取Kotlin Flow并将每个元素发送至服务
- 生产者:每5秒读取一次ArrayDeque,转换后发送至Flow
- 队列:ArrayDeque
相关代码:
private val consumerScope = CoroutineScope(dispatcherProvider.io() + SupervisorJob()) private val producerScope = CoroutineScope(dispatcherProvider.io() + SupervisorJob()) private val eventsQueueProducer: Flow<EventLogBatch> = channelFlow { while (currentCoroutineContext().isActive) { eventsQueue.get()?.let { log.error("Producer emitting logs to queue consumers") send(it) } delay(Delay) } } private val eventsQueueConsumer = eventsQueueProducer .onEach { event -> log.error("Received log events. Sending to backend") putLogEvents(requestBuilder.buildRequest( logsClient = it, logEvents = event.events.sortedBy { event -> event.timestamp } )) } override fun start() { eventsQueueConsumer.produceIn(consumerScope) eventsQueueProducer.launchIn(producerScope) } override suspend fun publishEvents(eventLogs: List<EventLog>) { eventLogs .onEach { log.error("Received event logs to chunk and send") } .map { it.toInputLogEvent() } .filter { it.message !== null } .chunkInMaxSizeBatches() .forEach { batch -> log.error("Enqueueing new log batch") eventsQueue.enqueueLogBatch(batch) } }
问题:已看到生产者日志Producer emitting logs to queue consumers,但消费者日志Received log events. Sending to backend数量不符,推测消费者工作异常。
可能的原因
1. 重复订阅Flow导致事件被“空消费”
在start()方法中同时执行了两个Flow订阅操作:
eventsQueueConsumer.produceIn(consumerScope) eventsQueueProducer.launchIn(producerScope)
channelFlow属于冷流,每次订阅都会启动独立的生产者循环。这里的两个调用相当于启动了两个并行的生产者:
- 第一个订阅绑定了消费者逻辑,会处理事件并打印消费者日志
- 第二个订阅仅单纯订阅生产者Flow,没有任何处理逻辑,相当于“空消费”了一部分事件,这部分事件不会触发消费者日志,最终导致消费者日志数量少于生产者。
2. ArrayDeque线程不安全导致数据丢失
ArrayDeque不是线程安全集合,当前代码中:
publishEvents在协程中调用enqueueLogBatch添加元素- 生产者在另一个协程中调用
get()取出元素
多协程并发操作时,可能出现队列状态不一致(比如扩容时的元素错乱、取出时的索引异常),导致部分元素被重复取出或直接丢失,造成生产者和消费者处理的事件数量不匹配。
3. 消费者协程因异常终止
如果putLogEvents调用AWS服务时抛出异常,onEach块中没有异常处理逻辑,会直接导致消费者协程崩溃。虽然使用了SupervisorJob,但崩溃的协程不会自动重启,后续生产者发送的事件将无人处理,造成消费者日志数量停滞。
4. eventsQueue.get()的逻辑缺陷
如果get()方法仅读取队列元素而不移除,生产者会在每次循环中重复发送同一个元素,导致生产者日志重复;如果get()的实现存在其他逻辑问题(比如只能取出最后一个元素、队列空时返回null的时机错误),也会导致事件处理不完整,最终日志数量不符。
内容的提问来源于stack exchange,提问作者colymore
相关产品推荐
相关产品推荐

