8000 Optimize ZPipeline.fromSink by adamgfraser · Pull Request #7988 · zio/zio · GitHub
[go: up one dir, main page]
More Web Proxy on the site http://driver.im/
Skip to content

Optimize ZPipeline.fromSink #7988

New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Merged
merged 1 commit into from
Apr 5, 2023
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
60 changes: 25 additions & 35 deletions streams/shared/src/main/scala/zio/stream/ZPipeline.scala
Original file line number Diff line number Diff line change
Expand Up @@ -1445,55 +1445,45 @@ object ZPipeline extends ZPipelinePlatformSpecificConstructors {
)(implicit trace: Trace): ZPipeline[Env, Err, In, Out] =
new ZPipeline(
ZChannel.suspend {
val leftovers: AtomicReference[Chunk[Chunk[In]]] = new AtomicReference(Chunk.empty)
val upstreamDone: AtomicBoolean = new AtomicBoolean(false)
var leftover: Chunk[In] = Chunk.empty
var upstreamDone: Boolean = false

lazy val buffer: ZChannel[Any, Err, Chunk[In], Any, Err, Chunk[In], Any] =
lazy val upstream: ZChannel[Any, Err, Chunk[In], Any, Err, Chunk[In], Any] =
ZChannel.suspend {
val l = leftovers.get
val l = leftover

if (l.isEmpty)
ZChannel.readWith(
(c: Chunk[In]) => ZChannel.write(c) *> buffer,
(e: Err) => ZChannel.fail(e),
(done: Any) => ZChannel.succeedNow(done)
ZChannel.readWithCause(
(c: Chunk[In]) => ZChannel.write(c) *> upstream,
(e: Cause[Err]) => ZChannel.failCause(e),
(done: Any) => {
upstreamDone = true
ZChannel.succeedNow(done)
}
)
else {
leftovers.set(Chunk.empty)
ZChannel.writeChunk(l) *> buffer
leftover = Chunk.empty
ZChannel.write(l) *> upstream
}
}

def concatAndGet(c: Chunk[Chunk[In]]): Chunk[Chunk[In]] = {
val ls = leftovers.get
val concat = ls ++ c.filter(_.nonEmpty)
leftovers.set(concat)
concat
}

lazy val upstreamMarker: ZChannel[Any, Err, Chunk[In], Any, Err, Chunk[In], Any] =
ZChannel.readWith(
(in: Chunk[In]) => ZChannel.write(in) *> upstreamMarker,
(err: Err) => ZChannel.fail(err),
(done: Any) => ZChannel.succeed(upstreamDone.set(true)) *> ZChannel.succeedNow(done)
lazy val writeDone: ZChannel[Any, Err, Chunk[In], Out, Err, Chunk[Out], Any] =
ZChannel.readWithCause(
elem => {
leftover ++= elem
writeDone
},
err => ZChannel.failCause(err),
out => ZChannel.write(Chunk.single(out))
)

lazy val transducer: ZChannel[Env, ZNothing, Chunk[In], Any, Err, Chunk[Out], Unit] =
sink.channel.collectElements.flatMap { case (leftover, z) =>
ZChannel
.succeed((upstreamDone.get, concatAndGet(leftover)))
.flatMap { case (done, newLeftovers) =>
val nextChannel =
if (done && newLeftovers.isEmpty) ZChannel.unit
else transducer

ZChannel.write(Chunk.single(z)) *> nextChannel
}
sink.channel.pipeTo(writeDone).flatMap { _ =>
if (upstreamDone && leftover.isEmpty) ZChannel.unit
else transducer
}

upstreamMarker >>>
buffer pipeToOrFail
transducer
upstream pipeToOrFail transducer
}
)

Expand Down
0