This is an automated email from the ASF dual-hosted git repository.
hepin pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-pekko.git
The following commit(s) were added to refs/heads/main by this push:
new 16e587ded1 perf: Fold In&OutHanlder in FlatMapPrefix.
16e587ded1 is described below
commit 16e587ded12d6ad83c66123c144b1dddfdebaad4
Author: He-Pin <[email protected]>
AuthorDate: Sun Dec 24 16:03:33 2023 +0800
perf: Fold In&OutHanlder in FlatMapPrefix.
---
.../pekko/stream/impl/fusing/FlatMapPrefix.scala | 47 +++++++++++-----------
1 file changed, 23 insertions(+), 24 deletions(-)
diff --git
a/stream/src/main/scala/org/apache/pekko/stream/impl/fusing/FlatMapPrefix.scala
b/stream/src/main/scala/org/apache/pekko/stream/impl/fusing/FlatMapPrefix.scala
index ce17848df1..9894870c63 100644
---
a/stream/src/main/scala/org/apache/pekko/stream/impl/fusing/FlatMapPrefix.scala
+++
b/stream/src/main/scala/org/apache/pekko/stream/impl/fusing/FlatMapPrefix.scala
@@ -131,38 +131,37 @@ import pekko.util.OptionVal
accumulated.clear()
subSource = OptionVal.Some(new
SubSourceOutlet[In]("FlatMapPrefix.subSource"))
val theSubSource = subSource.get
- theSubSource.setHandler {
- new OutHandler {
- override def onPull(): Unit = {
- if (!isClosed(in) && !hasBeenPulled(in)) {
- pull(in)
- }
- }
-
- override def onDownstreamFinish(cause: Throwable): Unit = {
- if (!isClosed(in)) {
- cancel(in, cause)
- }
- }
- }
- }
subSink = OptionVal.Some(new
SubSinkInlet[Out]("FlatMapPrefix.subSink"))
val theSubSink = subSink.get
- theSubSink.setHandler {
- new InHandler {
- override def onPush(): Unit = {
- push(out, theSubSink.grab())
- }
+ val handler = new InHandler with OutHandler {
+ override def onPush(): Unit = {
+ push(out, theSubSink.grab())
+ }
- override def onUpstreamFinish(): Unit = {
- complete(out)
+ override def onUpstreamFinish(): Unit = {
+ complete(out)
+ }
+
+ override def onUpstreamFailure(ex: Throwable): Unit = {
+ fail(out, ex)
+ }
+
+ override def onPull(): Unit = {
+ if (!isClosed(in) && !hasBeenPulled(in)) {
+ pull(in)
}
+ }
- override def onUpstreamFailure(ex: Throwable): Unit = {
- fail(out, ex)
+ override def onDownstreamFinish(cause: Throwable): Unit = {
+ if (!isClosed(in)) {
+ cancel(in, cause)
}
}
}
+
+ theSubSource.setHandler(handler)
+ theSubSink.setHandler(handler)
+
val matVal =
try {
val flow = f(prefix)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]