statefulMap 操作符详解:借助状态转换流的每个元素
statefulMap 操作符详解:借助状态转换流的每个元素
statefulMap 是 Akka Streams 中一个强大的流转换操作符,它允许你在处理流元素的同时维护一个状态,实现如 zipWithIndex、distinctUntilChanged、按条件分组等"有状态"的流处理逻辑。本文将从签名、语义、源码实现到完整示例,系统讲解如何在 Akka 项目中正确使用 statefulMap。
操作符定位与适用场景
statefulMap 属于 Akka Streams 的 简单操作符(Simple operators) 家族,其核心价值在于:将普通 map 的"一对一变换"升级为"携带状态的变换"。它特别适合以下场景:
- 为流中的每个元素附带自增索引(zipWithIndex 行为)
- 去重连续重复元素(distinctUntilChanged 行为)
- 按条件缓冲并成批下发元素
- 在流结束时基于累计状态产生最终输出(如分组剩余元素、汇总统计)
- 结合其他操作符实现更复杂的流处理(如模拟 statefulMapConcat)
如果你只需要无状态的元素变换,请使用 map,无需引入状态开销。
方法签名
statefulMap 在 Scala DSL 和 Java DSL 中均有对应 API,位于 Flow / Source / SubFlow / SubSource 上(如 scaladsl/Flow.scala):
Scala
def statefulMapS, T => S)(f: (S, Out) => (S, T), onComplete: S => Option[T]): Repr[T]
Java
def statefulMapS, T: Repr[T]
三个参数的职责:
| 参数 | 类型 | 说明 |
|---|---|---|
create |
() => S |
状态工厂函数,流被物化(materialized)时调用一次,返回初始状态,用于映射第一个元素 |
f |
(S, Out) => (S, T) |
映射函数,接收当前状态与上游元素,返回一对值:传给下一个映射函数的新状态 + 要向下游发射的元素 |
onComplete |
S => Option[T] |
完成函数,在流结束(上游完成/下游取消/流失败三者先到者)时调用一次,返回可选的最终输出元素 |
关于状态的类型,源码注释明确指出:映射函数返回的状态可以每次相同、可以是新的不可变状态,也允许使用可变状态(这在 Java 示例中体现得尤为明显,如直接复用 LinkedList/ArrayList 作为缓冲区)。
底层实现原理
statefulMap 在 Akka Streams 内部由 akka.stream.impl.fusing.Ops 包中的 StatefulMap GraphStage 实现(Ops.scala),它是一个标准的 GraphStage[FlowShape[In, Out]],理解其内部逻辑有助于把握操作符的精确语义。
状态生命周期
从源码可以看到状态的完整生命周期(Ops.scala):
override def preStart(): Unit = {
createNewState()
}
preStart() 阶段调用 create() 创建初始状态并保存。每个元素到达时(onPush),从 inlet 取出元素,调用映射函数 f(state.get, elem),将返回的新状态写回,同时把映射结果 push 到下游(Ops.scala)。
完成语义:onComplete 的三种触发时机
onComplete 函数在以下三种情况之一发生时恰好调用一次(对应文档中的"第一个到达者"语义):
- 上游正常完成(
onUpstreamFinish,Ops.scala):若onComplete返回Some(elem)且下游仍接受元素,则该元素在操作符完成前被发射;返回None则直接完成。 - 上游失败(
onUpstreamFailure/closeStateAndFail,Ops.scala):onComplete的返回值被忽略(completeStateIfNeeded的结果仅在"上游完成"分支用于发射)。 - 下游取消(
onDownstreamFinish,Ops.scala):同样忽略返回值。
该逻辑集中在 completeStateIfNeeded() 方法中(Ops.scala):
private def completeStateIfNeeded(): Option[Out] = {
state match {
case OptionVal.Some(s) =>
state = OptionVal.none[S]
onComplete(s)
case _ => None
}
}
注意其内部还通过 OptionVal 保证状态只被消费一次,并在 postStop() 中兜底调用(Ops.scala),确保资源清理路径完整。
状态非空约束
源码对状态有一个硬性约束(Ops.scala):
private def throwIfNoState(): Unit = {
if (state.isEmpty)
throw new NullStateException(
"State returned by stateFulMap create lambda or mapping function was null, which is not allowed. " +
"Use Option or Optional to represent presence of state if needed.")
}
create 或映射函数返回的状态不能为 null,否则抛出 NullStateException(该异常不会被监督策略覆盖,见 Ops.scala)。若确实需要表达"无状态",应使用 Option/Optional 包装——这正是下面示例中广泛采用 Option 的原因。
监督策略(SupervisionStrategy)
文档明确说明 statefulMap 遵循 ActorAttributes.SupervisionStrategy。源码中通过 inheritedAttributes.mandatoryAttribute<a href="https://link.gitcode.com/i/33a04009d09e44d83ef3b09f1fe51d6c" target="_blank">SupervisionStrategy].decider 获取决策器([Ops.scala),当映射函数抛出非致命异常时按策略处理(Ops.scala):
- Stop(默认):调用
closeStateAndFail(ex)结束流并传播失败,同时仍会尝试调用onComplete清理状态; - Resume:跳过当前元素,
pull(in)继续处理下一个元素,状态保持不变; - Restart:先尝试
completeStateIfNeeded()发射可能的最终元素,然后调用create()重建全新状态继续处理。
完整示例
以下四个示例均来自官方文档配套测试(StatefulMap.scala 与 StatefulMap.java),并已通过仓库中 FlowStatefulMapSpec.scala 的自动化测试验证。
示例一:实现 zipWithIndex(自增索引)
Scala
Source(List("A", "B", "C", "D"))
.statefulMap(() => 0L)((index, elem) => (index + 1, (elem, index)), _ => None)
.runForeach(println)
// prints
//(A,0)
//(B,1)
//(C,2)
//(D,3)
Java
Source.from(Arrays.asList("A", "B", "C", "D"))
.statefulMap(
() -> 0L,
(index, element) -> Pair.create(index + 1, Pair.create(element, index)),
indexOnComplete -> Optional.empty())
.runForeach(System.out::println, system);
// prints
// Pair(A,0)
// Pair(B,1)
// Pair(C,2)
// Pair(D,3)
状态就是 Long 类型的计数器:初始为 0,每次映射返回 (index + 1, (elem, index))——新状态是递增后的索引,发射元素是 (元素, 当前索引)。由于每个元素都独立发射,无需在完成时补发,onComplete 返回 None。
示例二:bufferUntilChanged(缓冲到元素变化再下发)
Scala
Source("A" :: "B" :: "B" :: "C" :: "C" :: "C" :: "D" :: Nil)
.statefulMap(() => List.empty[String])(
(buffer, element) =>
buffer match {
case head :: _ if head != element => (element :: Nil, buffer)
case _ => (element :: buffer, Nil)
},
buffer => Some(buffer))
.filter(_.nonEmpty)
.runForeach(println)
// prints
//List(A)
//List(B, B)
//List(C, C, C)
//List(D)
Java
Source.from(Arrays.asList("A", "B", "B", "C", "C", "C", "D"))
.statefulMap(
() -> (List<String>) new LinkedList<String>(),
(buffer, element) -> {
if (buffer.size() > 0 && (!buffer.get(0).equals(element))) {
return Pair.create(
new LinkedList<>(Collections.singletonList(element)),
Collections.unmodifiableList(buffer));
} else {
buffer.add(element);
return Pair.create(buffer, Collections.<String>emptyList());
}
},
Optional::ofNullable)
.filterNot(List::isEmpty)
.runForeach(System.out::println, system);
// prints
// [A]
// [B, B]
// [C, C, C]
// [D]
状态是元素缓冲区:当新元素与缓冲头部不同时,将已缓冲的列表整体发射、并以新元素重置缓冲(此时发射 Nil 表示无输出);相同时继续追加缓冲(发射 Nil)。onComplete 返回 Some(buffer) 把最后一段缓冲补发出去,再接 filter(_.nonEmpty) 丢弃中间过程的空列表。
示例三:distinctUntilChanged(去重连续重复)
Scala
Source("A" :: "B" :: "B" :: "C" :: "C" :: "C" :: "D" :: Nil)
.statefulMap(() => Option.empty[String])(
(lastElement, elem) =>
lastElement match {
case Some(head) if head == elem => (Some(elem), None)
case _ => (Some(elem), Some(elem))
},
_ => None)
.collect { case Some(elem) => elem }
.runForeach(println)
// prints
//A
//B
//C
//D
Java
Source.from(Arrays.asList("A", "B", "B", "C", "C", "C", "D"))
.statefulMap(
Optional::<String>empty,
(lastElement, element) -> {
if (lastElement.isPresent() && lastElement.get().equals(element)) {
return Pair.create(lastElement, Optional.<String>empty());
} else {
return Pair.create(Optional.of(element), Optional.of(element));
}
},
listOnComplete -> Optional.empty())
.via(Flow.flattenOptional())
.runForeach(System.out::println, system);
// prints
// A
// B
// C
// D
状态记录上一个元素(Option 类型):重复则发射 None(表示不输出),变化则发射 Some(elem)。随后用 collect(Scala)或 Flow.flattenOptional()(Java)过滤掉空输出。
示例四:结合 mapConcat 模拟 statefulMapConcat(每 3 个元素分组)
Scala
Source(1 to 10)
.statefulMap(() => List.empty[Int])(
(state, elem) => {
//grouped 3 elements into a list
val newState = elem :: state
if (newState.size == 3)
(Nil, newState.reverse)
else
(newState, Nil)
},
state => Some(state.reverse))
.mapConcat(identity)
.runForeach(println)
// prints
//1
//2
//3
//4
//5
//6
//7
//8
//9
//10
Java
Source.fromJavaStream(() -> IntStream.rangeClosed(1, 10))
.statefulMap(
() -> new ArrayList<Integer>(3),
(list, element) -> {
list.add(element);
if (list.size() == 3) {
return Pair.create(new ArrayList<Integer>(3), Collections.unmodifiableList(list));
} else {
return Pair.create(list, Collections.<Integer>emptyList());
}
},
Optional::ofNullable)
.mapConcat(list -> list)
.runForeach(System.out::println, system);
// prints
// 1
// 2
// 3
// 4
// 5
// 6
// 7
// 8
// 9
// 10
状态是累积缓冲区,攒满 3 个元素即以 newState.reverse 整组发射(倒序是因为 Scala 用 :: 头插),不满 3 个发射 Nil;onComplete 把不足一组的剩余元素 state.reverse 补发。输出经 mapConcat 摊平为单个元素。该模式可以等价实现 statefulMapConcat 的行为。
Reactive Streams 语义
按照 Reactive Streams 规范,statefulMap 的信号语义如下:
- emits(发射):当映射函数返回一个元素且下游准备好消费时
- backpressures(背压):当下游背压时
- completes(完成):当上游完成时
- cancels(取消):当下游取消时
这些语义与文档配套测试的断言一一对应(FlowStatefulMapSpec.scala),例如"happy case"测试验证了状态从 0 累加并逐个发射 (agg, elem) 后正常完成。
测试验证与典型行为
仓库的 FlowStatefulMapSpec.scala(共 406 行)系统覆盖了 statefulMap 的边界行为,可作为使用时的行为参考:
- happy case:基本累加 + 发射 + 完成(第 33-47 行)
- 完成时保留状态:最后不足分组的部分通过
onComplete补发(第 49-65 行) - Resume 监督:映射函数抛异常时跳过该元素继续处理(第 67 行起)
- Restart 监督:抛异常后重建状态继续处理
- 上游失败 / 下游取消:验证
onComplete在这些路径上的调用与返回值处理
使用注意事项小结
- 状态禁止为 null:
create与映射函数返回的状态都不能是null,需要表达空状态时用Option/Optional。 onComplete只调用一次:由上游完成、下游取消、流失败三者"先到者"触发;只有上游正常完成且下游仍可接收时,返回值才会被发射。- 状态可以可变:Java 示例直接复用
LinkedList/ArrayList作为可变状态是官方支持的做法,但需注意并发与一致性。 - 监督策略三态差异:Resume 保留状态跳过元素,Restart 重建状态,Stop 失败并清理;与普通 map 的监督行为有明显区别。
- 无状态需求请用 map:statefulMap 引入状态管理开销,纯变换场景应优先选择 map。