Akka Streaming:修复分块函数丢失最后数据块的问题
问题
我使用Akka Streaming向S3上传二进制文件,代码如下:
source .via(distributeChunks(MAX_BYTES_PER_CHUNK)) .throttle(maxRequestsPerSecond, 1.second, maximumBurst = 1, ThrottleMode.Shaping) .runFoldAsync(1) { (partNumber, bytes) => {...}
需要实现distributeChunks函数,将输入流分割为大小等于MAX_BYTES_PER_CHUNK的块,仅当最后剩余数据小于该值时生成更小的块。
我尝试了以下实现:
private def distributeChunks(maxChunkSize: Int): Flow[ByteString, ByteString, NotUsed] = Flow[ByteString] .statefulMapConcat { () => var buffer = ByteString.empty { bs: ByteString => buffer ++= bs val chunks = new ArrayBuffer[ByteString] while (buffer.length >= maxChunkSize) { val (chunk, rest) = buffer.splitAt(maxChunkSize) chunks += chunk buffer = rest } chunks.toList } } .mapMaterializedValue(_ => NotUsed)
该实现能生成符合大小要求的块,但会丢失最后一块数据,求优化代码实现需求。
测试用例:
- 文件大小10MB,最大允许块大小2MB:应分割为5个2MB块
- 文件大小9MB,最大允许块大小2MB:应分割为4个2MB块和1个1MB块
解决方案
问题根源是statefulMapConcat仅在每个输入元素到来时输出分割后的块,当流终止时,缓冲区里剩余的最后一块数据没有被输出。利用statefulMapConcat支持返回清理函数的特性,就能在流结束时输出剩余数据。
修复后的代码:
private def distributeChunks(maxChunkSize: Int): Flow[ByteString, ByteString, NotUsed] = Flow[ByteString] .statefulMapConcat { () => var buffer = ByteString.empty // 处理每个输入元素的逻辑 val processElement: ByteString => List[ByteString] = { bs => buffer ++= bs val chunks = new ArrayBuffer[ByteString] while (buffer.length >= maxChunkSize) { val (chunk, rest) = buffer.splitAt(maxChunkSize) chunks += chunk buffer = rest } chunks.toList } // 流终止时的清理逻辑:输出剩余缓冲区数据 val onComplete: () => List[ByteString] = () => if (buffer.nonEmpty) List(buffer) else Nil // 返回处理函数与清理函数的元组 processElement -> onComplete } .mapMaterializedValue(_ => NotUsed)
关键优化点
statefulMapConcat允许返回(处理函数, 清理函数)的元组,清理函数会在流结束时被调用。- 流终止时,清理函数检查缓冲区是否有剩余数据,若有则将其作为最后一个块输出,彻底解决最后一块丢失的问题。
- 原有分割逻辑保持不变,确保非最后一块的大小严格等于
maxChunkSize,最后一块小于该值时也能正常输出。
内容的提问来源于stack exchange,提问作者Shivam Sahil
相关产品推荐
相关产品推荐

