ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

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

statefulMap 操作符详解:借助状态转换流的每个元素 后端并发编程异步编程【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址https://gitcode.com/gh_mirrors/ak/akka-core点击查看免费下载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.scalaScaladef statefulMapS, T S)(f: (S, Out) (S, T), onComplete: S Option[T]): Repr[T]Javadef statefulMapS, T: Repr[T]三个参数的职责参数类型说明create() S状态工厂函数流被物化materialized时调用一次返回初始状态用于映射第一个元素f(S, Out) (S, T)映射函数接收当前状态与上游元素返回一对值传给下一个映射函数的新状态 要向下游发射的元素onCompleteS Option[T]完成函数在流结束上游完成/下游取消/流失败三者先到者时调用一次返回可选的最终输出元素关于状态的类型源码注释明确指出映射函数返回的状态可以每次相同、可以是新的不可变状态也允许使用可变状态这在 Java 示例中体现得尤为明显如直接复用LinkedList/ArrayList作为缓冲区。底层实现原理statefulMap 在 Akka Streams 内部由akka.stream.impl.fusing.Ops包中的StatefulMapGraphStage 实现Ops.scala它是一个标准的GraphStage[FlowShape[In, Out]]理解其内部逻辑有助于把握操作符的精确语义。状态生命周期从源码可以看到状态的完整生命周期Ops.scalaoverride def preStart(): Unit { createNewState() }preStart()阶段调用create()创建初始状态并保存。每个元素到达时onPush从 inlet 取出元素调用映射函数f(state.get, elem)将返回的新状态写回同时把映射结果 push 到下游Ops.scala。完成语义onComplete 的三种触发时机onComplete函数在以下三种情况之一发生时恰好调用一次对应文档中的第一个到达者语义上游正常完成onUpstreamFinishOps.scala若onComplete返回Some(elem)且下游仍接受元素则该元素在操作符完成前被发射返回None则直接完成。上游失败onUpstreamFailure/closeStateAndFailOps.scalaonComplete的返回值被忽略completeStateIfNeeded的结果仅在上游完成分支用于发射。下游取消onDownstreamFinishOps.scala同样忽略返回值。该逻辑集中在completeStateIfNeeded()方法中Ops.scalaprivate def completeStateIfNeeded(): Option[Out] { state match { case OptionVal.Some(s) state OptionVal.none[S] onComplete(s) case _ None } }注意其内部还通过OptionVal保证状态只被消费一次并在postStop()中兜底调用Ops.scala确保资源清理路径完整。状态非空约束源码对状态有一个硬性约束Ops.scalaprivate 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[SupervisionStrategy].decider获取决策器Ops.scala当映射函数抛出非致命异常时按策略处理Ops.scalaStop默认调用closeStateAndFail(ex)结束流并传播失败同时仍会尝试调用onComplete清理状态Resume跳过当前元素pull(in)继续处理下一个元素状态保持不变Restart先尝试completeStateIfNeeded()发射可能的最终元素然后调用create()重建全新状态继续处理。完整示例以下四个示例均来自官方文档配套测试StatefulMap.scala 与 StatefulMap.java并已通过仓库中 FlowStatefulMapSpec.scala 的自动化测试验证。示例一实现 zipWithIndex自增索引ScalaSource(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)JavaSource.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缓冲到元素变化再下发ScalaSource(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)JavaSource.from(Arrays.asList(A, B, B, C, C, C, D)) .statefulMap( () - (ListString) new LinkedListString(), (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.StringemptyList()); } }, 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去重连续重复ScalaSource(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 //DJavaSource.from(Arrays.asList(A, B, B, C, C, C, D)) .statefulMap( Optional::Stringempty, (lastElement, element) - { if (lastElement.isPresent() lastElement.get().equals(element)) { return Pair.create(lastElement, Optional.Stringempty()); } 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)。随后用collectScala或Flow.flattenOptional()Java过滤掉空输出。示例四结合 mapConcat 模拟 statefulMapConcat每 3 个元素分组ScalaSource(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 //10JavaSource.fromJavaStream(() - IntStream.rangeClosed(1, 10)) .statefulMap( () - new ArrayListInteger(3), (list, element) - { list.add(element); if (list.size() 3) { return Pair.create(new ArrayListInteger(3), Collections.unmodifiableList(list)); } else { return Pair.create(list, Collections.IntegeremptyList()); } }, 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 个发射NilonComplete把不足一组的剩余元素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在这些路径上的调用与返回值处理使用注意事项小结状态禁止为 nullcreate与映射函数返回的状态都不能是null需要表达空状态时用Option/Optional。onComplete只调用一次由上游完成、下游取消、流失败三者先到者触发只有上游正常完成且下游仍可接收时返回值才会被发射。状态可以可变Java 示例直接复用LinkedList/ArrayList作为可变状态是官方支持的做法但需注意并发与一致性。监督策略三态差异Resume 保留状态跳过元素Restart 重建状态Stop 失败并清理与普通 map 的监督行为有明显区别。无状态需求请用 mapstatefulMap 引入状态管理开销纯变换场景应优先选择 map。赞分享后端并发编程异步编程【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址https://gitcode.com/gh_mirrors/ak/akka-core点击查看免费下载相关推荐OpenCore Legacy Patcher终极指南让老Mac焕发新生的完整教程OpenCore Legacy Patcher终极指南让老Mac焕发新生的完整教程 OpenCore Legacy PatcherOCLP是一款革命性的开操作系统固件驱动开发RxJS 4 的 take 操作符详解从序列头部精确截取 N 个元素RxJS 4 的 take 操作符详解从序列头部精确截取 N 个元素 Rx.Observable.prototype.take count, schedule后端RxJS 4 操作符详解toArray —— 将 Observable 序列收敛为单个数组元素RxJS 4 操作符详解toArray —— 将 Observable 序列收敛为单个数组元素 导读 toArray 是 RxJSThe Reactive后端上一篇LaMa图像修复训练中断恢复指南掌握检查点与状态保存策略下一篇Brython与DevOps5个自动化构建、测试和部署流程的终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表