Skip to main content

Streams

zio.blocks.streams is a pull-based streaming library for Scala 3 (and Scala 2.13) with synchronous and asynchronous readers, typed errors, resource safety, and primitive specialization. Streams are lazy descriptions -- nothing executes until a terminal operation is driven. Cross-platform terminals ending in Async return Async[Either[E, Z]]; the JVM also provides blocking terminals returning Either[E, Z]. The library has zero runtime dependencies beyond ZIO Blocks modules, and avoids boxing on primitive element types through JVM-type-specialized internal readers.

ZIO Blocks Streams is built on three composable primitives:

TypeDescriptionKey operation
Stream[+E, +A]A lazy, pull-based sequence of elements that may fail with error Estream.via(pipe)
Pipeline[-In, +Out]A reusable, composable stream-to-stream transformationpipe.andThen(other)
Sink[+E, -A, +Z]A stream consumer that produces a typed result Zstream.run(sink)

Overview​

ZIO Blocks Streams is designed around three core principles:

Dual execution. Cross-platform *Async terminals drive either synchronous or asynchronous sources without blocking and require no ZIO runtime. On the JVM, plain terminals such as run, runCollect, and head are blocking compatibility twins.

Pull-based evaluation. Execution is driven from the consumer (Sink) backward through the pipeline to the source (Stream). This enables natural short-circuiting: if a sink only needs the first three elements, the stream stops producing after three elements — no work is wasted.

Resource safety via RAII. Resources acquired during stream construction (file handles, database connections, etc.) are always released in finally blocks, whether the stream succeeds, fails, or is short-circuited.

Quick Start​

Here's a minimal JVM example. Streams are lazy descriptions — nothing executes until you call a terminal operation like runCollect. Use runCollectAsync for the cross-platform form.

Unless a section says otherwise, examples using plain terminals (run, runCollect, head, and their peers) are JVM-only shorthand. Replace them with the matching *Async terminal in shared JVM/Scala.js code.

import zio.blocks.streams.*
import zio.blocks.chunk.Chunk

// Build a lazy stream description
val stream = Stream.range(1, 100)
.filter(_ % 2 == 0)
.map(_ * 3)

// Run it — nothing executes until here
val result = stream.take(5).runCollect
// Right(Chunk(6, 12, 18, 24, 30))

Installation​

Add the Streams module to your SBT build:

libraryDependencies += "dev.zio" %% "zio-blocks-streams" % "0.0.56"

For Scala.js (JavaScript/Node.js):

libraryDependencies += "dev.zio" %%% "zio-blocks-streams" % "0.0.56"

Supported Scala versions: 2.13.x and 3.x.

Why Streams?​

Streaming libraries in the Scala ecosystem typically require an effect system. fs2 runs in a cats.effect-compatible F[_], Kyo Streams needs the Kyo runtime, and Pekko (formerly Akka) Streams needs the actor runtime. When your code is synchronous and you want streaming without pulling in an effect monad, the options narrow considerably.

zio.blocks.streams fills that gap:

FeatureZB Streamsfs2KyoOxPekko
Effect system requiredNoYes (cats-effect)Yes (Kyo)No (virtual threads)Yes (Pekko)
Execution modelSync/async, pull-basedAsync, pull-basedAsync, chunkSynchronous, pull-basedAsync, push
Typed errorsEither[E, Z]Not verified hereNot verified hereNot verified hereNot verified here
Primitive specializationYes (zero boxing)Not verified hereNot verified hereNot verified hereNot verified here
Stack-safe deep pipelinesYes (trampolined)Not verified hereNot verified hereNot verified hereNot verified here
Resource safetyScope integrationResource/bracketKyo resourcestry/finallyGraph lifecycle
Dependenciesscope, chunk, combinators, ringbuffer, asyncfs2-corekyo-prelude, kyo-coreox corepekko-stream

The provider columns name the artifacts and versions this repository pins for benchmarking — fs2 3.14.0, Pekko 1.7.0, Kyo 1.0.0-RC6, Ox 1.0.6 (build.sbt, streams-benchmark/benchmark-manifest.json). Nothing outside the ZB Streams column is measured or verified in this repository.

Benchmarks​

The repository carries a JMH benchmark suite under streams-benchmark. Provider comparisons are governed by the allowlist and provider versions recorded in streams-benchmark/benchmark-manifest.json; results are only comparable within a single benchmark class and contract, and are not aggregated into a ranking here. Re-run the benchmarks on your target environment before drawing any performance conclusion.

If you are evaluating Scala 2 compatibility work, read the Scala 2 compatibility design note before moving any Stream or Sink hot-path combinators behind version-specific seams.

Core Mental Model​

To understand ZIO Blocks Streams fully, it's helpful to see how the three primitives fit together and how data flows through a pipeline from source to sink. This section walks through the architecture and explains each component in depth.

Execution Flow​

Operations on streams transform the pipeline and ultimately run it against a sink:

┌──────────────────────────────────┐
│ Stream[E, A] │
│ (lazy description) │
└──────────────────┬───────────────┘
│
.flatMap, .map, .filter, etc.
│
┌──────────────────▼───────────────┐
│ Pipeline[-In, +Out] │
│ (stream → stream transformation) │
└──────────────────┬───────────────┘
│
.via(pipe)
│
┌──────────────────▼───────────────┐
│ Sink[E, A, Z] │
│ (stream consumer → result Z) │
└──────────────────┬───────────────┘
│
.run(sink)
│
┌──────────────────▼───────────────┐
│ Async[Either[E, Z]] │
│ (or blocking Either on JVM) │
└──────────────────────────────────┘

The last box is where the flow forks. The same Stream description is materialized either as a synchronous reader, drained on the calling thread by a blocking terminal such as run or runCollect, or as an asynchronous reader, driven without blocking by the matching *Async terminal. Which engine runs is decided by the source and operators the pipeline is built from, not by the terminal you call. See Async Execution for how that classification works.

1) Stream[E, A] -- a Lazy Sequence​

A Stream[+E, +A] is a description of a potentially infinite sequence of elements of type A that may fail with an error of type E. It is covariant in both type parameters.

Nothing happens when you construct a stream or chain transformations. Execution only begins when you drive a terminal operation. Cross-platform terminals (runAsync, runCollectAsync, headAsync, countAsync, etc.) return Async[Either[E, Z]]; plain blocking terminals return Either[E, Z] on the JVM:

  • Left(e) -- a typed stream error
  • Right(z) -- the successful result

Untyped defects (unexpected exceptions) propagate as thrown exceptions, not as Left values.

import zio.blocks.streams.*

// This does nothing -- it's just a description
val description: Stream[Nothing, Int] =
Stream.range(0, 1_000_000)
.filter(_ % 7 == 0)
.map(_ * 2)
.take(100)

// Only this line executes the pipeline
// val result = description.runCollect

Streams render their pipeline structure as a human-readable string:

val s = Stream.range(0, 100).map(_ + 1).filter(_ > 50).take(10)
println(s) // Stream.range(0, 100).map(...).filter(...).take(10)

This makes debugging and logging straightforward -- you can see exactly what transformations a stream applies without running it.


2) Sink[E, A, Z] -- a Consumer​

A Sink[+E, -A, +Z] consumes elements of type A from a stream and produces a final result of type Z. Sinks are passed to Stream.run:

import zio.blocks.streams.*

val streamSinks = Stream.range(1, 101)

// Built-in sinks
val total = streamSinks.run(Sink.count)
val items = streamSinks.run(Sink.collectAll)
val sum = streamSinks.run(Sink.sumInt)
val first = streamSinks.run(Sink.head)

Most sinks also have convenience methods directly on Stream:

stream.count // Either[Nothing, Long]
stream.runCollect // Either[Nothing, Chunk[Int]]
stream.head // Either[Nothing, Option[Int]]
stream.last // Either[Nothing, Option[Int]]

Sinks compose with contramap (pre-process input) and map (post-process result):

val lengthSink: Sink[Nothing, String, Long] =
Sink.sumInt.contramap[Int, String](_.length)

val doubled: Sink[Nothing, Int, Long] =
Sink.sumInt.map(_ * 2)

3) Pipeline[In, Out] -- Reusable Transformation​

A Pipeline[-In, +Out] is a reusable stream transformation. It decouples the transformation logic from any specific stream, so you can define it once and apply it many times.

// Define a reusable pipeline
val normalize: Pipeline[Int, Double] =
Pipeline.filter[Int](_ > 0)
.andThen(Pipeline.map[Int, Double](_.toDouble / 100.0))

// Apply to different streams
val result1 = Stream.range(-10, 10).via(normalize).runCollect
val result2 = Stream.fromIterable(List(42, -5, 100, 0)).via(normalize).runCollect

Pipelines compose with andThen:

val step1: Pipeline[String, Int] =
Pipeline.map[String, Int](_.length)

val step2: Pipeline[Int, Int] =
Pipeline.filter[Int](_ > 3)

val combined: Pipeline[String, Int] =
step1.andThen(step2)

You can also apply a pipeline to a sink with andThenSink / applyToSink, which pre-processes the sink's input:

val countLong: Sink[Nothing, String, Long] =
Pipeline.map[String, Int](_.length)
.andThenSink(Sink.sumInt)

Synchronous and Asynchronous Execution​

There is one Stream type. It serves both execution modes, and there is no mode type parameter, no AsyncStream, and no annotation to write.

The type that decides is Reader, not Stream. Materializing a stream yields either a Reader.SyncReader[A], whose pulls return values directly, or a Reader.AsyncReader[A], whose pulls return Async values. A pipeline that is synchronous end to end materializes as the former; a single asynchronous source or operator anywhere in it makes the whole pipeline asynchronous.

The *Async terminals -- runAsync, runCollectAsync, runDrainAsync, runFoldAsync, countAsync, headAsync, and their peers -- are the cross-platform API. They all return Async[Either[E, Z]]: the typed error E stays inside the Either, while the outer Async fails only on a defect. They drive a synchronous pipeline just as correctly as an asynchronous one, so shared JVM/Scala.js code can use them unconditionally.

The blocking terminals -- run, runCollect, runDrain, runFold, count, head, and their peers -- are JVM-only compatibility twins returning a bare Either[E, Z]. They do not exist on Scala.js, and cross-compiled sources cannot call them.

  • Async Execution -- the full execution model: classification, the Reader union, the async source constructors, operators, and terminals.
  • Platform Differences -- the availability matrix of every member that differs between the JVM and Scala.js.

Error Handling​

Streams distinguish between two kinds of failures:

  • Typed errors (E) — domain errors you expect and handle, returned as Left in the result. Use catchAll, orElse, or mapError to recover.
  • Defects (Throwable) — unexpected exceptions from bugs or system failures. Use catchDefect to recover, or they propagate as thrown exceptions.
val failing: Stream[String, Int] =
Stream.fromIterable(List(1, 2, 3)) ++ Stream.fail("oops") ++ Stream.fromIterable(List(4, 5))

val recovered = failing.catchAll(_ => Stream.fromIterable(List(99)))
recovered.runCollect // Right(Chunk(1, 2, 3, 99))

Resource Management​

Streams integrate with zio.blocks.scope.Scope for deterministic resource cleanup. The fromAcquireRelease constructor guarantees that a resource is acquired lazily (when the stream runs), used to produce elements, and then released — even if the stream is short-circuited early via take(), fails with an error, or succeeds normally. The release function is wired into a finally block, ensuring cleanup always happens.

import zio.blocks.streams.*

val managed = Stream.fromAcquireRelease(
acquire = scala.io.Source.fromFile("data.txt"),
release = _.close()
)(source => Stream.fromIterator(source.getLines()))

managed.take(10).runCollect
// File is closed in finally block regardless of outcome

This eliminates the need for manual try/finally when working with resources — the stream handles it for you.

Primitive Specialization​

ZB Streams carries the JVM representation of the element type through the whole pipeline, so a stream of primitives is not boxed at each stage boundary. Specialization is not limited to Int, Long, Float, and Double: there are nine logical lanes -- the eight primitive pull identities Boolean, Byte, Short, Char, Int, Long, Float, and Double, plus the reference fallback -- and the synchronous interpreter compacts them into five storage lanes: int-like (Boolean/Byte/Short/Char/Int), Long, Float, Double, and reference. Nine logical lanes therefore does not mean nine interpreter arrays; the operator tag selects the identity-specific reads over the shared storage.

Zero-Boxing Streams explains how a lane is chosen and what the specialization evidence is for.

import zio.blocks.streams.*

// This entire pipeline runs with ZERO boxing of the Int elements.
// Every step uses specialized readInt/writeInt internally.
val sum: Either[Nothing, Long] =
Stream.range(0, 1_000_000) // Int-specialized source
.filter(_ % 2 == 0) // Int-specialized filter
.map(_ * 3) // Int->Int specialized map
.runFold(0L)(_ + _) // Long-specialized accumulator

This matters most for numeric workloads — data processing, statistics, encoding/decoding — where millions of elements flow through multi-stage pipelines.

Practical Guidance​

  • Start with Stream constructors and terminal operations. You can get very far with Stream.range, Stream.fromIterable, .map, .filter, and .runCollect.
  • Use Either pattern matching to handle the result: Right(value) for success, Left(error) for typed failures.
  • Prefer Stream.fromAcquireRelease when wrapping resources (files, connections, etc.) over manual try/finally. It guarantees cleanup even on early termination via take, head, or error.
  • Use the auto-closing I/O constructors (fromInputStream, fromJavaReader, NioStreams.fromChannel) by default. Only use the Unmanaged variants when you need to borrow a resource whose lifetime is managed elsewhere.
  • Use Pipeline when you have a transformation you want to reuse across multiple streams or apply to sinks.
  • Use && for zipping instead of manual zip calls. Tuples flatten automatically: a && b && c produces (A, B, C) not ((A, B), C).
  • Leverage primitive specialization for numeric workloads. Streams of any primitive element type avoid boxing automatically on the synchronous path; use Sink.sumInt, runFold(0)(_ + _), etc. for allocation-free folds.
  • Use scan for running accumulators, grouped for batching, and sliding for windowed computations.
  • Use render/toString to inspect pipeline structure during debugging — it shows each transformation stage without executing the stream.
  • Use Sink.create (JVM only) as an escape hatch when none of the built-in sinks fit; on Scala.js and in cross-compiled code use Sink.createAsync or Sink.createBoth.
  • suspend is your friend for recursive or self-referential stream definitions, preventing stack overflow during construction.
  • Typed errors vs. defects: use Stream.fail for expected domain errors and Stream.die for programmer errors. Use catchAll for the former, catchDefect for the latter.

Usage Examples​

This section shows practical examples of using streams in real-world scenarios. Each subsection demonstrates a different aspect of the API with runnable code examples.

Creating Streams​

Here are the most common ways to construct a stream. Choose the constructor that best fits your data source:

import zio.blocks.streams.*
import zio.blocks.chunk.Chunk

// From explicit elements
Stream.fromIterable(List(1, 2, 3)) // Stream[Nothing, Int]
Stream.fromIterable(List("a", "b", "c")) // Stream[Nothing, String]

// From collections
Stream.fromChunk(Chunk(1, 2, 3)) // Stream[Nothing, Int]
Stream.fromIterable(List("x", "y", "z")) // Stream[Nothing, String]
Stream.fromIterator(Iterator.from(1)) // Stream[Nothing, Int] (lazy)

// Ranges
Stream.range(0, 100) // 0 to 99
Stream.fromRange(1 to 50) // 1 to 50

// Single values (primitive-specialized)
Stream.succeed(42) // Stream[Nothing, Int]
Stream.succeed(3.14) // Stream[Nothing, Double]
Stream.succeed("hello") // Stream[Nothing, String]

// Special streams
Stream.empty // Stream[Nothing, Nothing]
Stream.fail("error") // Stream[String, Nothing]
// Stream.die(new Exception("defect")) // throws on evaluation

// Generators
Stream.repeat(1) // infinite stream of 1s
Stream.unfold(0)(n => // 0, 1, 2, ..., 9
if n < 10 then Some((n, n + 1)) else None
)

// Side-effects
Stream.eval(println("hello")) // prints, emits nothing
Stream.attempt(someFallibleCall()) // captures exceptions as typed errors
Stream.attemptEval(riskyEffect()) // same, for Unit-returning effects

// Deferred construction (useful for recursion)
Stream.suspend(expensiveStreamBuilder())

// I/O sources (auto-closing) - JVM only
Stream.fromInputStream(inputStream) // Stream[IOException, Byte] (auto-closes)
Stream.fromJavaReader(javaReader) // Stream[IOException, Char] (auto-closes)

// I/O sources (borrowing -- caller manages lifetime) - JVM only
Stream.fromInputStreamUnmanaged(inputStream) // Stream[IOException, Byte] (does NOT close)
Stream.fromJavaReaderUnmanaged(javaReader) // Stream[IOException, Char] (does NOT close)

Transforming Streams​

Streams support many transformation operations. Use map for element-wise changes, filter for selection, and flatMap for expanding elements into sub-streams. See the Stream reference page for comprehensive examples of all transformation methods including map, filter, flatMap, collect, scan, mapAccum, distinct, intersperse, and more.


Zipping Streams with &&​

The && operator zips two streams element-by-element into tuples. The resulting stream ends when either input is exhausted.

import zio.blocks.streams.*

val names: Stream[Nothing, String] = Stream.fromIterable(List("Alice", "Bob", "Charlie"))
val ages: Stream[Nothing, Int] = Stream.fromIterable(List(30, 25, 35))
val ids: Stream[Nothing, Long] = Stream.fromIterable(List(1L, 2L, 3L))

// Two-way zip
val pairs = names && ages
pairs.runCollect // Right(Chunk(("Alice", 30), ("Bob", 25), ("Charlie", 35)))

// Three-way zip -- tuples flatten automatically
val triples = names && ages && ids
triples.runCollect // Right(Chunk(("Alice", 30, 1L), ("Bob", 25, 2L), ("Charlie", 35, 3L)))

When the error types differ, they widen through Concat — to a union E1 | E2 on Scala 3, and to a meaningful least upper bound (or Either[E1, E2] for disjoint types) on Scala 2.13:

import zio.blocks.streams.*

sealed trait MyError
val s1: Stream[MyError, Int] = Stream.fromIterable(List(1, 2, 3))
sealed trait OtherError
val s2: Stream[OtherError, Int] = Stream.fromIterable(List(4, 5, 6))
// val zipped = s1 && s2

Primitive Specialization​

All nine logical lanes are specialized, not only Int, Long, Float, and Double. Every intermediate step uses the identity-specific read (readInt, readByte, readChar, and so on), so no java.lang.Integer wrappers are allocated between stages.

// This entire pipeline runs with ZERO boxing of the Int elements.
// Every step uses specialized readInt/writeInt internally.
val sum: Either[Nothing, Long] =
Stream.range(0, 1_000_000) // Int-specialized source
.filter(_ % 2 == 0) // Int-specialized filter
.map(_ * 3) // Int->Int specialized map
.runFold(0L)(_ + _) // Long-specialized accumulator

This matters most for numeric workloads -- data processing, statistics, encoding/decoding -- where millions of elements flow through multi-stage pipelines.


Consuming Streams​

Terminal operations run the stream and produce a final result. Use runCollect to gather all elements, runDrain to discard them, or specialized operations like head, count, and foldLeft:

val s = Stream.range(1, 11) // 1 to 10

// Collect all elements
s.runCollect // Right(Chunk(1, 2, 3, ..., 10))

// Discard all elements (run for side-effects only)
s.tapEach(println).runDrain

// Fold
s.runFold(0)(_ + _) // Right(55) (Int accumulator)
s.runFold(0L)(_ + _) // Right(55L) (Long accumulator)
s.runFold(0.0)(_ + _) // Right(55.0) (Double accumulator)

// Foreach
s.runForeach(n => println(n))
s.foreach(n => println(n)) // alias

// Aggregates
s.count // Right(10L)
s.head // Right(Some(1))
s.last // Right(Some(10))
s.exists(_ > 5) // Right(true)
s.forall(_ > 0) // Right(true)
s.find(_ > 7) // Right(Some(8))

// Run with an explicit Sink
s.run(Sink.sumInt) // Right(55L)
s.run(Sink.take(3)) // Right(Chunk(1, 2, 3))

Async Execution​

runCollectAsync is the cross-platform twin of runCollect. It returns Async[Either[E, Chunk[A]]] rather than Either[E, Chunk[A]], so it never blocks and compiles on both the JVM and Scala.js:

import zio.blocks.streams._
import zio.blocks.chunk.Chunk
import zio.blocks.async._

val evens: Stream[Nothing, Int] =
Stream.range(1, 100).filter(_ % 2 == 0).map(_ * 3)

// Still a description -- nothing has run.
val pending: Async[Either[Nothing, Chunk[Int]]] = evens.take(5).runCollectAsync

// Stay in Async: map the result rather than extracting it.
val described: Async[String] = pending.map {
case Right(values) => s"collected ${values.length} elements"
case Left(error) => s"failed: $error"
}

Something has to drive the Async at the edge of the world. On the JVM that is .block, which belongs in main or a test and nowhere else:

import zio.blocks.streams._
import zio.blocks.chunk.Chunk
import zio.blocks.async._

// JVM only: Async#block throws on Scala.js.
val result: Either[Nothing, Chunk[Int]] =
Stream.range(1, 100).filter(_ % 2 == 0).take(5).runCollectAsync.block
// Right(Chunk(2, 4, 6, 8, 10))

Scala.js code keeps the Async and hands it to the host runtime instead. See Async Execution for the full terminal family and Platform Differences for what is available where.


Error Handling Patterns​

Streams support two types of failures: typed errors that you can handle explicitly, and defects (exceptions) that propagate. Here are common patterns for dealing with both:

// Typed error: appears in Either
val result = Stream.fail("not found").runCollect
// result: Left("not found")

// Recover and continue
val safe =
Stream.fromIterable(List(1, 2)) ++ Stream.fail("oops") ++ Stream.fromIterable(List(3))
val recovered = safe.catchAll(_ => Stream.fromIterable(List(99))).runCollect
// Right(Chunk(1, 2, 99))

// Transform error type by catching and converting
val inputError: Stream[String, Int] = Stream.fail("bad input")
val transformed = inputError.catchAll(msg => Stream.fail(new IllegalArgumentException(msg)))

// Fallback stream
val primary: Stream[String, Int] = Stream.fail("down")
val backup: Stream[String, Int] = Stream.fromIterable(List(1, 2, 3))
val result2 = (primary || backup).runCollect
// Right(Chunk(1, 2, 3))

// Catch defects (unexpected exceptions)
val risky: Stream[Nothing, Int] =
Stream.fromIterable(List(1, 2, 3)).map { n =>
if n == 2 then throw new ArithmeticException("boom")
else n
}

val handled = risky.catchDefect {
case _: ArithmeticException => Stream.fromIterable(List(0))
}.runCollect
// Right(Chunk(1, 0))

Resource Safety Patterns​

When working with files, network connections, or other resources, use the resource-safe constructors to guarantee cleanup. Here are the most common patterns:

import zio.blocks.streams.*
import zio.blocks.scope.*

// Bracket pattern: acquire/use/release
def fileLines(path: String): Stream[Nothing, String] =
Stream.fromAcquireRelease(
acquire = scala.io.Source.fromFile(path),
release = _.close()
) { source =>
Stream.fromIterable(source.getLines().toList)
}

// Compose resource-safe streams -- both resources are released
val merged =
fileLines("input1.txt") ++ fileLines("input2.txt")

// Only reads 10 lines; both files are still closed properly
merged.take(10).runCollect

// ensuring: attach a finalizer
var cleaned = false
Stream.range(1, 6)
.ensuring { cleaned = true }
.take(2)
.runDrain
// cleaned == true, even though only 2 of 5 elements were consumed

// defer: register cleanup that runs on stream close
val withDefer =
Stream.defer(println("releasing lock")) ++
Stream.range(1, 100)

NIO Integration (JVM Only)​

On the JVM, NioStreams and NioSinks provide zero-copy integration with java.nio buffers and channels.

NioStreams -- Creating Streams From NIO Sources​

import zio.blocks.streams.*
import java.nio.ByteBuffer
import java.nio.channels.FileChannel
import java.nio.file.{Paths, StandardOpenOption}

// From a ByteBuffer
val buf = ByteBuffer.wrap(Array[Byte](1, 2, 3, 4, 5))
NioStreams.fromByteBuffer(buf).runCollect
// Right(Chunk(1, 2, 3, 4, 5))

// Typed buffer views (zero-boxing)
val intBuf = ByteBuffer.allocate(16).putInt(1).putInt(2).putInt(3).putInt(4).flip()
NioStreams.fromByteBufferInt(intBuf).runCollect
// Right(Chunk(1, 2, 3, 4))

// Similarly: fromByteBufferLong, fromByteBufferFloat, fromByteBufferDouble

// From a ReadableByteChannel (auto-closing)
val ch = FileChannel.open(Paths.get("data.bin"), StandardOpenOption.READ)
val bytes = NioStreams.fromChannel(ch, bufSize = 4096).runCollect
// ch is closed automatically when the stream completes

// From a ReadableByteChannel (borrowing -- caller manages lifetime)
val ch2 = FileChannel.open(Paths.get("data.bin"), StandardOpenOption.READ)
val bytes2 = NioStreams.fromChannelUnmanaged(ch2, bufSize = 4096).runCollect
ch2.close() // caller is responsible for closing

NioSinks -- Writing to NIO Targets​

import zio.blocks.streams.*
import zio.blocks.chunk.Chunk
import java.nio.ByteBuffer
import java.nio.channels.FileChannel
import java.nio.file.{Files, StandardOpenOption}

// Write to a ByteBuffer using a typed sink (Int values, zero-boxing)
val outBuf = ByteBuffer.allocate(1024)
Stream.range(1, 5).run(NioSinks.fromByteBufferInt(outBuf))

// Write to a WritableByteChannel using a stream of bytes
val tempPath = Files.createTempFile("zio-blocks-streams-", ".bin")
val outCh = FileChannel.open(
tempPath,
StandardOpenOption.WRITE,
StandardOpenOption.TRUNCATE_EXISTING
)
val bytes = Chunk.fromIterable(List[Byte](1, 2, 3, 4, 5))
try Stream.fromChunk(bytes).run(NioSinks.fromChannel(outCh))
finally {
outCh.close()
Files.deleteIfExists(tempPath)
}

Pipeline Composition​

Pipelines are composable transformations that can be reused across different streams. Build complex transformations by chaining pipelines together with andThen:

import zio.blocks.streams.*

// Build reusable transformation steps
val parseInts: Pipeline[String, Int] =
Pipeline.collect[String, Int] {
case s if s.matches("-?\\d+") => s.toInt
}

val positiveOnly: Pipeline[Int, Int] =
Pipeline.filter[Int](_ > 0)

val doubled: Pipeline[Int, Int] =
Pipeline.map[Int, Int](_ * 2)

// Compose into a single pipeline
val fullPipeline: Pipeline[String, Int] =
parseInts
.andThen(positiveOnly)
.andThen(doubled)

// Apply to any stream of strings
Stream.fromIterable(List("10", "abc", "-3", "7", "0", "25"))
.via(fullPipeline)
.runCollect
// Right(Chunk(20, 14, 50))

// Apply to a sink (pre-process the sink's input)
val sumPositiveDoubled: Sink[Nothing, String, Long] =
fullPipeline.andThenSink(Sink.sumInt)

Stream.fromIterable(List("10", "abc", "-3", "7", "0", "25"))
.run(sumPositiveDoubled)
// Right(84L)

See Also​

  • Async Execution -- the synchronous/asynchronous execution model, the Reader union, and the *Async terminal family
  • Platform Differences -- which members exist on the JVM, on Scala.js, and on both
  • Zero-Boxing Streams -- the primitive lanes and how one is chosen
  • Async -- the Async effect type the cross-platform terminals return
  • Mux -- coordinating many keyed streams over one shared transport