Pen*_*gin 4 scala akka akka-stream
学习 Akka 流。我有一个记录流,每个时间单位有很多记录,已经按时间排序(来自 Slick),并且我想将它们分批放入时间组中,以便通过检测时间步长何时发生变化来进行处理。
例子
case class Record(time: Int, payload: String)
Run Code Online (Sandbox Code Playgroud)
如果传入的流是
Record(1, "a")
Record(1, "k")
Record(1, "k")
Record(1, "a")
Record(2, "r")
Record(2, "o")
Record(2, "c")
Record(2, "k")
Record(2, "s")
Record(3, "!")
...
Run Code Online (Sandbox Code Playgroud)
我想把它变成
Batch(1, Seq("a","k","k","a"))
Batch(2, Seq("r","o","c","k","s"))
Batch(3, Seq("!"))
...
Run Code Online (Sandbox Code Playgroud)
到目前为止,我只发现按固定数量的记录进行分组,或者分成许多子流,但从我的角度来看,我不需要多个子流。
更新:我发现batch,但它看起来更关心背压,而不仅仅是一直批处理。
statefulMapConcat是 Akka Streams 库中的多功能工具。
val records =
Source(List(
Record(1, "a"),
Record(1, "k"),
Record(1, "k"),
Record(1, "a"),
Record(2, "r"),
Record(2, "o"),
Record(2, "c"),
Record(2, "k"),
Record(2, "s"),
Record(3, "!")
))
.concat(Source.single(Record(0, "notused"))) // needed to print the last element
records
.statefulMapConcat { () =>
var currentTime = 0
var payloads: Seq[String] = Nil
record =>
if (record.time == currentTime) {
payloads = payloads :+ record.payload
Nil
} else {
val previousState = (currentTime, payloads)
currentTime = record.time
payloads = Seq(record.payload)
List(previousState)
}
}
.runForeach(println)
Run Code Online (Sandbox Code Playgroud)
运行上面的命令会打印以下内容:
(0,List())
(1,List(a, k, k, a))
(2,List(r, o, c, k, s))
(3,List(!))
Run Code Online (Sandbox Code Playgroud)
您可以调整示例以打印Batch对象。
| 归档时间: |
|
| 查看次数: |
1230 次 |
| 最近记录: |