statefulMap 操作符详解:借助状态转换流的每个元素

原创2026-09-22 17:32:08160 阅读
文章标签:后端并发编程异步编程

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 函数在以下三种情况之一发生时恰好调用一次(对应文档中的"第一个到达者"语义):

  1. 上游正常完成(onUpstreamFinish,Ops.scala):若 onComplete 返回 Some(elem) 且下游仍接受元素,则该元素在操作符完成前被发射;返回 None 则直接完成。
  2. 上游失败(onUpstreamFailure / closeStateAndFail,Ops.scala):onComplete 的返回值被忽略(completeStateIfNeeded 的结果仅在"上游完成"分支用于发射)。
  3. 下游取消(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 在这些路径上的调用与返回值处理

使用注意事项小结

  1. 状态禁止为 null:create 与映射函数返回的状态都不能是 null,需要表达空状态时用 Option/Optional。
  2. onComplete 只调用一次:由上游完成、下游取消、流失败三者"先到者"触发;只有上游正常完成且下游仍可接收时,返回值才会被发射。
  3. 状态可以可变:Java 示例直接复用 LinkedList/ArrayList 作为可变状态是官方支持的做法,但需注意并发与一致性。
  4. 监督策略三态差异:Resume 保留状态跳过元素,Restart 重建状态,Stop 失败并清理;与普通 map 的监督行为有明显区别。
  5. 无状态需求请用 map:statefulMap 引入状态管理开销,纯变换场景应优先选择 map。
登录后查看全文
akka-core