我经常使用纯函数式、非 akka 技术来解决此类问题,然后将这些函数“提升”到 akka 结构中。我的意思是我尝试只使用 scala “东西”,然后稍后尝试将这些东西包装在 akka 中......
文件创建
从基于“随机生成的名称”的FileOutputStream 创建开始:
def randomFileNameGenerator : String = ??? //not specified in question
import java.io.FileOutputStream
val randomFileOutGenerator : () => FileOutputStream =
() => new FileOutputStream(randomFileNameGenerator)
州
需要某种方式来存储当前文件的“状态”,例如已写入的字节数:
case class FileState(byteCount : Int = 0,
fileOut : FileOutputStream = randomFileOutGenerator())
文件写入
首先,我们确定是否会超出给定ByteString 的最大文件大小阈值:
import akka.util.ByteString
val isEndOfChunk : (FileState, ByteString, Int) => Boolean =
(state, byteString, maxBytes) =>
state.byteCount + byteString.length > maxBytes
然后我们必须编写一个函数来创建一个新的FileState,如果我们已经用尽了当前的容量,或者如果它仍然低于容量,则返回当前状态:
val closeFileInState : FileState => Unit =
(_ : FileState).fileOut.close()
val getCurrentFileState(FileState, ByteString, Int) => FileState =
(state, byteString, maxBytes) =>
if(isEndOfChunk(maxBytes, state, byteString)) {
closeFileInState(state)
FileState()
}
else
state
剩下的就是写信给FileOutputStream:
val writeToFileAndReturn(FileState, ByteString) => FileState =
(fileState, byteString) => {
fileState.fileOut write byteString.toArray
fileState copy (byteCount = fileState.byteCount + byteString.size)
}
//the signature ordering will become useful
def writeToChunkedFile(maxBytes : Int)(fileState : FileState, byteString : ByteString) : FileState =
writeToFileAndReturn(getCurrentFileState(maxBytes, fileState, byteString), byteString)
折叠任何 GenTraversableOnce
在 Scala 中,GenTraversableOnce 是任何具有折叠运算符的集合,无论是否并行。其中包括 Iterator, Vector, Array, Seq, scala stream, ...最终的writeToChunkedFile函数与GenTraversableOnce#fold的签名完美匹配:
val anyIterable : Iterable = ???
val finalFileState = anyIterable.fold(FileState())(writetochunkedFile(maxBytes))
最后一个松散的结局;最后一个FileOutputStream 也需要关闭。由于折叠只会发出最后一个FileState,我们可以关闭那个:
closeFileInState(finalFileState)
Akka 流
Akka Flow 的 fold 来自 FlowOps#fold,恰好与 GenTraversableOnce 签名匹配。因此,我们可以将常规函数“提升”为流值,类似于我们使用 Iterable fold 的方式:
import akka.stream.scaladsl.Flow
def chunkerFlow(maxBytes : Int) : Flow[ByteString, FileState, _] =
Flow[ByteString].fold(FileState())(writeToChunkedFile(maxBytes))
使用常规函数处理问题的好处在于,它们可以在流之外的其他异步框架中使用,例如期货或演员。在单元测试中,您也不需要 akka ActorSystem,只需要常规语言数据结构。
import akka.stream.scaladsl.Sink
import scala.concurrent.Future
def byteStringSink(maxBytes : Int) : Sink[ByteString, _] =
chunkerFlow(maxBytes) to (Sink foreach closeFileInState)
然后你可以使用这个Sink 来消耗来自HttpRequest 的HttpEntity。