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

如何用Kotlin Coroutines实现异步预拉取批次数据的迭代器

基于Kotlin协程的预取迭代器实现

核心实现逻辑用协程Deferred存储下一页的预取结果,替代原线程方案的Thread实例,拉取操作默认运行在Dispatchers.IO调度器适配IO密集型的外部服务调用,完全兼容原返回Sequence的接口要求,调用方无需修改现有遍历代码。

完整实现代码

import kotlinx.coroutines.Deferred
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.async
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.cancel

fun readStructuredLogs(...): Sequence<Payload.JsonPayload> {
    return object : Iterator<LogPayloadType>, AutoCloseable {
        private val logging: Logging = options.service
        private var currentPage: Page<LogEntry> = logging.listLogEntries(Logging.EntryListOption.filter(buildLogFilter(...)))
        private var i = 0
        private var batch: List<LogEntry> = currentPage.values.toList()
        private var isClosed = false
        // 存储下一页的预取任务
        private var nextPageDeferred: Deferred<Page<LogEntry>?>? = null

        override fun close() {
            // 取消还在运行的预取任务避免资源泄漏
            nextPageDeferred?.cancel()
            logging.close()
            isClosed = true
        }

        override fun hasNext(): Boolean {
            if (isClosed) return false
            val hasNext = i < batch.size || currentPage.hasNextPage()
            if (!hasNext) close()
            return hasNext
        }

        override fun next(): LogPayloadType {
            if (!hasNext()) throw NoSuchElementException()

            // 当前批次遍历到倒数第二个元素时启动下一页预取,最大化并行时间
            if (i == batch.size - 1 && currentPage.hasNextPage() && nextPageDeferred == null) {
                nextPageDeferred = async(Dispatchers.IO) {
                    currentPage.nextPage
                }
            }

            // 当前批次已经遍历完,等待预取结果
            if (i == batch.size) {
                currentPage = runBlocking { nextPageDeferred!!.await() } ?: throw NoSuchElementException()
                batch = currentPage.values.toList()
                i = 0
                // 拿到当前页后立刻启动下下页的预取
                nextPageDeferred = if (currentPage.hasNextPage()) {
                    async(Dispatchers.IO) { currentPage.nextPage }
                } else null
            }

            val logEntry = batch[i++]
            return logEntry.getPayload()
        }
    }.asSequence()
}

流程匹配说明

  • 首次拉取第1个批次同步完成,迭代器返回第一个批次倒数第二个元素时就会启动第2个批次的后台预取
  • 调用方处理第1个批次剩余元素的过程中,协程后台并行拉取第2个批次
  • 第1个批次处理完成时,第2个批次的预取大概率已经完成,直接拿到结果无需等待
  • 拿到第2个批次的同时自动启动第3个批次的预取,以此类推

注意事项

  • 如果外部服务拉取逻辑是阻塞调用,必须指定Dispatchers.IO调度器,避免阻塞协程调度
  • 调用方完全不需要修改原有的Sequence遍历逻辑,直接使用forEach、map等操作即可
  • 相比于线程实现版本,协程版本复用IO调度器线程池,资源消耗更低,高并发场景下性能优势更明显

内容的提问来源于stack exchange,提问作者Lutosław

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 22:06:03