object Subscriber extends Serializable
- Source
- Subscriber.scala
- Alphabetic
- By Inheritance
- Subscriber
- Serializable
- Serializable
- AnyRef
- Any
- Hide All
- Show All
- Public
- All
Type Members
-
implicit final
class
Extensions
[T] extends AnyVal
Extension methods for Subscriber.
-
trait
Sync
[-T] extends Subscriber[T] with Observer.Sync[T]
A
Subscriber.Sync
is a Subscriber whoseonNext
signal is synchronous (i.e.A
Subscriber.Sync
is a Subscriber whoseonNext
signal is synchronous (i.e. the upstream observable doesn't need to wait on aFuture
in order to decide whether to send the next event or not).
Value Members
-
def
apply[T](observer: Observer[T], scheduler: Scheduler): Subscriber[T]
Subscriber builder
-
def
canceled[A](implicit s: Scheduler): Sync[A]
Helper for building an empty subscriber that doesn't do anything, but that returns
Stop
ononNext
. -
def
dump[A](prefix: String, out: PrintStream = System.out)(implicit s: Scheduler): Sync[A]
Builds an Subscriber that just logs incoming events.
-
def
empty[A](implicit s: Scheduler): Sync[A]
Helper for building an empty subscriber that doesn't do anything, besides logging errors in case they happen.
-
def
fromReactiveSubscriber[T](subscriber: org.reactivestreams.Subscriber[T], subscription: Cancelable)(implicit s: Scheduler): Subscriber[T]
Given an
org.reactivestreams.Subscriber
as defined by the Reactive Streams specification, it builds an Subscriber instance compliant with the Monix Rx implementation. -
def
toReactiveSubscriber[T](source: Subscriber[T], bufferSize: Int): org.reactivestreams.Subscriber[T]
Transforms the source Subscriber into a
org.reactivestreams.Subscriber
instance as defined by the Reactive Streams specification.Transforms the source Subscriber into a
org.reactivestreams.Subscriber
instance as defined by the Reactive Streams specification.- bufferSize
a strictly positive number, representing the size of the buffer used and the number of elements requested on each cycle when communicating demand, compliant with the reactive streams specification
-
def
toReactiveSubscriber[T](subscriber: Subscriber[T]): org.reactivestreams.Subscriber[T]
Transforms the source Subscriber into a
org.reactivestreams.Subscriber
instance as defined by the Reactive Streams specification. - object Sync extends Serializable
This is the API documentation for the Monix library.
Package Overview
monix.execution exposes lower level primitives for dealing with asynchronous execution:
Atomic
types, as alternative tojava.util.concurrent.atomic
monix.eval is for dealing with evaluation of results, thus exposing Task and Coeval.
monix.reactive exposes the
Observable
pattern:Observable
implementationsmonix.types implements type-class shims, to be translated to type-classes provided by libraries such as Cats or Scalaz.
monix.cats is the optional integration with the Cats library, providing translations for the types described in
monix.types
.monix.scalaz is the optional integration with the Scalaz library, providing translations for the types described in
monix.types
.