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
相关产品推荐
相关产品推荐

