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

Kotlin中如何并行执行S3文件读取与解析操作?

问题描述

需要异步调用两次S3读取并解析数据,两个文件总记录约24000条。尝试用Kotlin协程async实现,但两个操作串行执行(文件1在10:29:53启动,文件2在10:30:53才启动),希望两个操作能同时启动。

当前代码如下:

fun process() {
    runBlocking {
        async { parseFile("File1") }.await()
        async { parseFile("File2") }.await()

    }
}

private fun parseFile(fileName: String) {
    log.info { "Initiate read from S3 for $fileName " }
    val getRecord = s3RepositoryClient.getObjectContentInputStream(fileName)
    val parseToString=parse(
        objectMapper, getRecord,"test"
    )
    parseToString.parallelStream().forEach {
        dynamoDbClient.save(it)
    }
}


fun parseToString(mapper: ObjectMapper, records: InputStream, textName: String): List<Billing> {
    var response: List<Billing> = emptyList()
    val reader = BufferedReader(records.reader())
    try {
        var line = reader.readLine()
        while (line != null) {
            response = response.plus(
                Billing.from(
                    mapper.readValue(line, Employee::class.java),
                    line,
                    textName
                )
            )
            line = reader.readLine()
        }
    } finally {
        reader.close()
    }
    return response
}

问题原因

当前代码中,你先调用async { parseFile("File1") }.await(),这会等待第一个异步任务完全执行完成后,才会创建并启动第二个async任务,导致两个操作串行执行,出现时间间隔。

解决方案

1. 同时启动异步任务再统一等待

先创建两个异步任务对象,再统一等待它们完成,这样两个任务会同时启动:

fun process() {
    runBlocking {
        // 先创建两个异步任务,此时它们会立即开始执行
        val job1 = async { parseFile("File1") }
        val job2 = async { parseFile("File2") }
        
        // 等待两个任务都完成,也可以用awaitAll(job1, job2)简化
        job1.await()
        job2.await()
    }
}

2. 将阻塞操作切换到IO线程池

parseFile中的S3读取、文件解析、DynamoDB写入都是阻塞IO操作,直接在协程默认线程执行会阻塞线程,影响并发效率。建议用Dispatchers.IO专门处理这类操作:

private suspend fun parseFile(fileName: String) = withContext(Dispatchers.IO) {
    log.info { "Initiate read from S3 for $fileName " }
    val getRecord = s3RepositoryClient.getObjectContentInputStream(fileName)
    val parsedList = parse(objectMapper, getRecord, "test")
    
    parsedList.forEach {
        dynamoDbClient.save(it)
    }
}

这里把parseFile改为suspend函数,通过withContext(Dispatchers.IO)将阻塞操作切换到IO线程池,避免占用协程计算线程。

3. 优化解析函数性能

原parseToString用response.plus()每次创建新列表,会产生大量临时对象,改用mutableListOf添加元素更高效:

fun parseToString(mapper: ObjectMapper, records: InputStream, textName: String): List<Billing> {
    val response = mutableListOf<Billing>()
    val reader = BufferedReader(records.reader())
    try {
        var line = reader.readLine()
        while (line != null) {
            response.add(
                Billing.from(
                    mapper.readValue(line, Employee::class.java),
                    line,
                    textName
                )
            )
            line = reader.readLine()
        }
    } finally {
        reader.close()
    }
    return response
}

最终优化代码示例

import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.async
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withContext
import kotlinx.coroutines.awaitAll

fun process() {
    runBlocking {
        val job1 = async { parseFile("File1") }
        val job2 = async { parseFile("File2") }
        awaitAll(job1, job2)
    }
}

private suspend fun parseFile(fileName: String) = withContext(Dispatchers.IO) {
    log.info { "Initiate read from S3 for $fileName " }
    // 用use自动关闭流,避免资源泄漏
    s3RepositoryClient.getObjectContentInputStream(fileName).use { inputStream ->
        val parsedList = parse(objectMapper, inputStream, "test")
        parsedList.forEach { dynamoDbClient.save(it) }
    }
}

fun parse(mapper: ObjectMapper, records: InputStream, textName: String): List<Billing> {
    val response = mutableListOf<Billing>()
    BufferedReader(records.reader()).use { reader ->
        var line = reader.readLine()
        while (line != null) {
            response.add(
                Billing.from(
                    mapper.readValue(line, Employee::class.java),
                    line,
                    textName
                )
            )
            line = reader.readLine()
        }
    }
    return response
}

内容的提问来源于stack exchange,提问作者rahul.cs

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.04 02:23:25