Package io.micronaut.http.body.stream
Class PublisherAsBlocking
java.lang.Object
io.micronaut.http.body.stream.PublisherAsBlocking
- All Implemented Interfaces:
Closeable,AutoCloseable,org.reactivestreams.Subscriber<io.micronaut.core.io.buffer.ReadBuffer>
@Internal
public final class PublisherAsBlocking
extends Object
implements org.reactivestreams.Subscriber<io.micronaut.core.io.buffer.ReadBuffer>, Closeable
A subscriber that allows blocking reads from a publisher. Handles resource cleanup properly.
- Since:
- 4.2.0
-
Constructor Summary
Constructors -
Method Summary
Modifier and TypeMethodDescriptionvoidclose()@Nullable ThrowableThe failure fromonError(Throwable).voidvoidvoidonNext(io.micronaut.core.io.buffer.ReadBuffer o) voidonSubscribe(org.reactivestreams.Subscription s) @Nullable io.micronaut.core.io.buffer.ReadBuffertake()Get the next object.
-
Constructor Details
-
PublisherAsBlocking
public PublisherAsBlocking()
-
-
Method Details
-
getFailure
The failure fromonError(Throwable). Whentake()returnsnull, this may be set if the reactive stream ended in failure.- Returns:
- The failure, or
nullif either the stream is not done, or the stream completed successfully.
-
onSubscribe
public void onSubscribe(org.reactivestreams.Subscription s) - Specified by:
onSubscribein interfaceorg.reactivestreams.Subscriber<io.micronaut.core.io.buffer.ReadBuffer>
-
onNext
public void onNext(io.micronaut.core.io.buffer.ReadBuffer o) - Specified by:
onNextin interfaceorg.reactivestreams.Subscriber<io.micronaut.core.io.buffer.ReadBuffer>
-
onError
- Specified by:
onErrorin interfaceorg.reactivestreams.Subscriber<io.micronaut.core.io.buffer.ReadBuffer>
-
onComplete
public void onComplete()- Specified by:
onCompletein interfaceorg.reactivestreams.Subscriber<io.micronaut.core.io.buffer.ReadBuffer>
-
take
@Nullable public @Nullable io.micronaut.core.io.buffer.ReadBuffer take() throws InterruptedExceptionGet the next object.- Returns:
- The next object, or
nullif the stream is done - Throws:
InterruptedException
-
close
public void close()- Specified by:
closein interfaceAutoCloseable- Specified by:
closein interfaceCloseable
-