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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 13:35:30