Class Pipe<T>

java.lang.Object
io.vertx.mutiny.core.streams.Pipe<T>
All Implemented Interfaces:
MutinyDelegate

public class Pipe<T> extends Object implements MutinyDelegate
Pipe data from a ReadStream to a WriteStream and performs flow control where necessary to prevent the write stream buffer from getting overfull.

Instances of this class read items from a ReadStream and write them to a WriteStream. If data can be read faster than it can be written this could result in the write queue of the WriteStream growing without bound, eventually causing it to exhaust all available RAM.

To prevent this, after each write, instances of this class check whether the write queue of the WriteStream is full, and if so, the ReadStream is paused, and a drainHandler is set on the WriteStream.

When the WriteStream has processed half of its backlog, the drainHandler will be called, which results in the pump resuming the ReadStream.

This class can be used to pipe from any ReadStream to any WriteStream, e.g. from an HttpServerRequest to an AsyncFile, or from NetSocket to a WebSocket.

Please see the documentation for more information.

NOTE: This class has been automatically generated from the original non Mutiny-ified interface.

See Also:
  • Pipe
  • Field Details

    • __TYPE_ARG

      public static final TypeArg<Pipe> __TYPE_ARG
    • __typeArg_0

      public final TypeArg<T> __typeArg_0
  • Constructor Details

    • Pipe

      public Pipe(io.vertx.core.streams.Pipe<T> delegate)
      Create a new instance of Pipe delegating to the given (non-null) instance of Pipe.
    • Pipe

      public Pipe(io.vertx.core.streams.Pipe<T> delegate, TypeArg<T> typeArg_0)
    • Pipe

      public Pipe(Object delegate, TypeArg<T> typeArg_0)
  • Method Details

    • getDelegate

      public io.vertx.core.streams.Pipe<T> getDelegate()
      Get the delegate instance.

      This method returns the instance on which this shim is delegating the calls. And so, give you access to the bare API.

      Specified by:
      getDelegate in interface MutinyDelegate
      Returns:
      the delegate instance
    • to

      @CheckReturnValue public io.smallrye.mutiny.Uni<Void> to(WriteStream<T> dst)
      Start to pipe the elements to the destination WriteStream.

      When the operation fails with a write error, the source stream is resumed.

      Unlike the bare Vert.x variant, this method returns a Uni. The uni emits the result of the operation as item. If the operation fails, the uni emits the failure.

      Don't forget to subscribe on it to trigger the operation.

      Parameters:
      dst - the destination write stream
      Returns:
      A Uni representing the asynchronous result of this operation.
      See Also:
      • io.vertx.core.streams.Pipe#to(WriteStream)
    • toAndAwait

      public void toAndAwait(WriteStream<T> dst)
      Start to pipe the elements to the destination WriteStream.

      When the operation fails with a write error, the source stream is resumed.

      Unlike the bare Vert.x variant, this method returns a Void. This method awaits indefinitely for the completion of the underlying asynchronous operation. If the operation completes successfully, the result is returned, otherwise the failure is thrown (potentially wrapped in a RuntimeException).

      Parameters:
      dst - the destination write stream
      See Also:
      • io.vertx.core.streams.Pipe#to(WriteStream)
    • toAndForget

      public Pipe<T> toAndForget(WriteStream<T> dst)
      Start to pipe the elements to the destination WriteStream.

      When the operation fails with a write error, the source stream is resumed.

      Unlike the bare Vert.x variant, this method ignores the Void result or any failure.

      Parameters:
      dst - the destination write stream
      Returns:
      The current instance to chain operations if needed.
      See Also:
      • io.vertx.core.streams.Pipe#to(WriteStream)
    • endOnFailure

      public Pipe<T> endOnFailure(boolean end)
      Set to true to call WriteStream.end() when the source ReadStream fails, false otherwise.
      Parameters:
      end - true to end the stream on a source ReadStream failure
      Returns:
      a reference to this, so the API can be used fluently
    • endOnSuccess

      public Pipe<T> endOnSuccess(boolean end)
      Set to true to call WriteStream.end() when the source ReadStream succeeds, false otherwise.
      Parameters:
      end - true to end the stream on a source ReadStream success
      Returns:
      a reference to this, so the API can be used fluently
    • endOnComplete

      public Pipe<T> endOnComplete(boolean end)
      Set to true to call WriteStream.end() when the source ReadStream completes, false otherwise.

      Calling this overwrites endOnFailure(boolean) and endOnSuccess(boolean).

      Parameters:
      end - true to end the stream on a source ReadStream completion
      Returns:
      a reference to this, so the API can be used fluently
    • close

      public void close()
      Close the pipe.

      The streams handlers will be unset and the read stream resumed unless it is already ended.

    • newInstance

      public static <T> Pipe<T> newInstance(io.vertx.core.streams.Pipe<T> delegate)
      Creates a new instance of the Pipe.
    • newInstance

      public static <T> Pipe<T> newInstance(io.vertx.core.streams.Pipe<T> delegate, TypeArg<T> typeArg_0)
      Creates a new instance of the Pipe.
    • hashCode

      public int hashCode()
      Overrides:
      hashCode in class Object
    • equals

      public boolean equals(Object o)
      Overrides:
      equals in class Object
    • toString

      public String toString()
      Overrides:
      toString in class Object