o

akka.stream.javadsl

StreamConverters

object StreamConverters

Converters for interacting with the blocking java.io streams APIs and Java 8 Streams

Source
StreamConverters.scala
Linear Supertypes
Ordering
  1. Alphabetic
  2. By Inheritance
Inherited
  1. StreamConverters
  2. AnyRef
  3. Any
  1. Hide All
  2. Show All
Visibility
  1. Public
  2. All

Value Members

  1. final def !=(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  2. final def ##(): Int
    Definition Classes
    AnyRef → Any
  3. final def ==(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  4. def asInputStream(readTimeout: FiniteDuration): Sink[ByteString, InputStream]

    Creates a Sink which when materialized will return an java.io.InputStream which it is possible to read the values produced by the stream this Sink is attached to.

    Creates a Sink which when materialized will return an java.io.InputStream which it is possible to read the values produced by the stream this Sink is attached to.

    This Sink is intended for inter-operation with legacy APIs since it is inherently blocking.

    You can configure the default dispatcher for this Source by changing the akka.stream.blocking-io-dispatcher or set it for a given Source by using akka.stream.ActorAttributes.

    The InputStream will be closed when the stream flowing into this Sink completes, and closing the InputStream will cancel this Sink.

    readTimeout

    the max time the read operation on the materialized InputStream should block

  5. def asInputStream(): Sink[ByteString, InputStream]

    Creates a Sink which when materialized will return an java.io.InputStream which it is possible to read the values produced by the stream this Sink is attached to.

    Creates a Sink which when materialized will return an java.io.InputStream which it is possible to read the values produced by the stream this Sink is attached to.

    This method uses a default read timeout, use #inputStream(FiniteDuration) to explicitly configure the timeout.

    This Sink is intended for inter-operation with legacy APIs since it is inherently blocking.

    You can configure the default dispatcher for this Source by changing the akka.stream.blocking-io-dispatcher or set it for a given Source by using akka.stream.ActorAttributes.

    The InputStream will be closed when the stream flowing into this Sink completes, and closing the InputStream will cancel this Sink.

  6. final def asInstanceOf[T0]: T0
    Definition Classes
    Any
  7. def asJavaStream[T](): Sink[T, Stream[T]]

    Creates a sink which materializes into Java 8 Stream that can be run to trigger demand through the sink.

    Creates a sink which materializes into Java 8 Stream that can be run to trigger demand through the sink. Elements emitted through the stream will be available for reading through the Java 8 Stream.

    The Java 8 Stream will be ended when the stream flowing into this Sink completes, and closing the Java Stream will cancel the inflow of this Sink.

    Java 8 Stream throws exception in case reactive stream failed.

    Be aware that Java Stream blocks current thread while waiting on next element from downstream. As it is interacting wit blocking API the implementation runs on a separate dispatcher configured through the akka.stream.blocking-io-dispatcher.

  8. def asOutputStream(): Source[ByteString, OutputStream]

    Creates a Source which when materialized will return an java.io.OutputStream which it is possible to write the ByteStrings to the stream this Source is attached to.

    Creates a Source which when materialized will return an java.io.OutputStream which it is possible to write the ByteStrings to the stream this Source is attached to. The write timeout for OutputStreams materialized will default to 5 seconds, @see #outputStream(FiniteDuration) if you want to override it.

    This Source is intended for inter-operation with legacy APIs since it is inherently blocking.

    You can configure the default dispatcher for this Source by changing the akka.stream.blocking-io-dispatcher or set it for a given Source by using akka.stream.ActorAttributes.

    The created OutputStream will be closed when the Source is cancelled, and closing the OutputStream will complete this Source.

  9. def asOutputStream(writeTimeout: FiniteDuration): Source[ByteString, OutputStream]

    Creates a Source which when materialized will return an java.io.OutputStream which it is possible to write the ByteStrings to the stream this Source is attached to.

    Creates a Source which when materialized will return an java.io.OutputStream which it is possible to write the ByteStrings to the stream this Source is attached to.

    This Source is intended for inter-operation with legacy APIs since it is inherently blocking.

    You can configure the default dispatcher for this Source by changing the akka.stream.blocking-io-dispatcher or set it for a given Source by using akka.stream.ActorAttributes.

    The created OutputStream will be closed when the Source is cancelled, and closing the OutputStream will complete this Source.

    writeTimeout

    the max time the write operation on the materialized OutputStream should block

  10. def clone(): AnyRef
    Attributes
    protected[java.lang]
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  11. final def eq(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  12. def equals(arg0: Any): Boolean
    Definition Classes
    AnyRef → Any
  13. def finalize(): Unit
    Attributes
    protected[java.lang]
    Definition Classes
    AnyRef
    Annotations
    @throws( classOf[java.lang.Throwable] )
  14. def fromInputStream(in: Creator[InputStream]): Source[ByteString, CompletionStage[IOResult]]

    Creates a Source from an java.io.InputStream created by the given function.

    Creates a Source from an java.io.InputStream created by the given function. Emitted elements are ByteString elements, chunked by default by 8192 bytes, except the last element, which will be up to 8192 in size.

    You can configure the default dispatcher for this Source by changing the akka.stream.blocking-io-dispatcher or set it for a given Source by using akka.stream.ActorAttributes.

    It materializes a CompletionStage of IOResult containing the number of bytes read from the source file upon completion, and a possible exception if IO operation was not completed successfully.

    The created InputStream will be closed when the Source is cancelled.

  15. def fromInputStream(in: Creator[InputStream], chunkSize: Int): Source[ByteString, CompletionStage[IOResult]]

    Creates a Source from an java.io.InputStream created by the given function.

    Creates a Source from an java.io.InputStream created by the given function. Emitted elements are chunkSize sized akka.util.ByteString elements, except the final element, which will be up to chunkSize in size.

    You can configure the default dispatcher for this Source by changing the akka.stream.blocking-io-dispatcher or set it for a given Source by using akka.stream.ActorAttributes.

    It materializes a CompletionStage containing the number of bytes read from the source file upon completion.

    The created InputStream will be closed when the Source is cancelled.

  16. def fromJavaStream[O, S <: BaseStream[O, S]](stream: Creator[BaseStream[O, S]]): Source[O, NotUsed]

    Creates a source that wraps a Java 8 Stream.

    Creates a source that wraps a Java 8 Stream. Source uses a stream iterator to get all its elements and send them downstream on demand.

    Example usage: Source.fromJavaStream(() -> IntStream.rangeClosed(1, 10))

    You can use Source.async to create asynchronous boundaries between synchronous java stream and the rest of flow.

  17. def fromOutputStream(f: Creator[OutputStream], autoFlush: Boolean): Sink[ByteString, CompletionStage[IOResult]]

    Sink which writes incoming ByteStrings to an OutputStream created by the given function.

    Sink which writes incoming ByteStrings to an OutputStream created by the given function.

    Materializes a CompletionStage of IOResult that will be completed with the size of the file (in bytes) at the streams completion, and a possible exception if IO operation was not completed successfully.

    You can configure the default dispatcher for this Source by changing the akka.stream.blocking-io-dispatcher or set it for a given Source by using akka.stream.ActorAttributes.

    The OutputStream will be closed when the stream flowing into this Sink is completed. The Sink will cancel the stream when the OutputStream is no longer writable.

    f

    A Creator which creates an OutputStream to write to

    autoFlush

    If true the OutputStream will be flushed whenever a byte array is written

  18. def fromOutputStream(f: Creator[OutputStream]): Sink[ByteString, CompletionStage[IOResult]]

    Sink which writes incoming ByteStrings to an OutputStream created by the given function.

    Sink which writes incoming ByteStrings to an OutputStream created by the given function.

    Materializes a CompletionStage of IOResult that will be completed with the size of the file (in bytes) at the streams completion, and a possible exception if IO operation was not completed successfully.

    You can configure the default dispatcher for this Source by changing the akka.stream.blocking-io-dispatcher or set it for a given Source by using akka.stream.ActorAttributes.

    This method uses no auto flush for the java.io.OutputStream @see Boolean) if you want to override it.

    The OutputStream will be closed when the stream flowing into this Sink is completed. The Sink will cancel the stream when the OutputStream is no longer writable.

    f

    A Creator which creates an OutputStream to write to

  19. final def getClass(): Class[_]
    Definition Classes
    AnyRef → Any
  20. def hashCode(): Int
    Definition Classes
    AnyRef → Any
  21. final def isInstanceOf[T0]: Boolean
    Definition Classes
    Any
  22. def javaCollector[T, R](collector: Creator[Collector[T, _, R]]): Sink[T, CompletionStage[R]]

    Creates a sink which materializes into a CompletionStage which will be completed with a result of the Java 8 Collector transformation and reduction operations.

    Creates a sink which materializes into a CompletionStage which will be completed with a result of the Java 8 Collector transformation and reduction operations. This allows usage of Java 8 streams transformations for reactive streams. The Collector will trigger demand downstream. Elements emitted through the stream will be accumulated into a mutable result container, optionally transformed into a final representation after all input elements have been processed. The Collector can also do reduction at the end. Reduction processing is performed sequentially

    Note that a flow can be materialized multiple times, so the function producing the Collector must be able to handle multiple invocations.

  23. def javaCollectorParallelUnordered[T, R](parallelism: Int)(collector: Creator[Collector[T, _, R]]): Sink[T, CompletionStage[R]]

    Creates a sink which materializes into a CompletionStage which will be completed with a result of the Java 8 Collector transformation and reduction operations.

    Creates a sink which materializes into a CompletionStage which will be completed with a result of the Java 8 Collector transformation and reduction operations. This allows usage of Java 8 streams transformations for reactive streams. The Collector will trigger demand downstream. Elements emitted through the stream will be accumulated into a mutable result container, optionally transformed into a final representation after all input elements have been processed. Collector can also do reduction at the end. Reduction processing is performed in parallel based on graph Balance.

    Note that a flow can be materialized multiple times, so the function producing the Collector must be able to handle multiple invocations.

  24. final def ne(arg0: AnyRef): Boolean
    Definition Classes
    AnyRef
  25. final def notify(): Unit
    Definition Classes
    AnyRef
  26. final def notifyAll(): Unit
    Definition Classes
    AnyRef
  27. final def synchronized[T0](arg0: ⇒ T0): T0
    Definition Classes
    AnyRef
  28. def toString(): String
    Definition Classes
    AnyRef → Any
  29. final def wait(): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  30. final def wait(arg0: Long, arg1: Int): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )
  31. final def wait(arg0: Long): Unit
    Definition Classes
    AnyRef
    Annotations
    @throws( ... )

Inherited from AnyRef

Inherited from Any

Ungrouped