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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 13:27:51