public class ZeroMQChannelStream<IN,OUT> extends ChannelStream<IN,OUT>
ReactorChannel.ConsumerSpec
log
Constructor and Description |
---|
ZeroMQChannelStream(Environment env,
long prefetch,
Dispatcher eventsDispatcher,
InetSocketAddress remoteAddress,
Codec<Buffer,IN,OUT> codec) |
Modifier and Type | Method and Description |
---|---|
void |
close() |
org.zeromq.ZMQ.Socket |
delegate() |
void |
doDecoded(IN in) |
protected void |
doSubscribeWriter(org.reactivestreams.Publisher<? extends OUT> writer,
org.reactivestreams.Subscriber<? super Void> postWriter) |
ReactorChannel.ConsumerSpec |
on()
Assign event handlers to certain channel lifecycle events.
|
InetSocketAddress |
remoteAddress()
Get the address of the remote peer.
|
ZeroMQChannelStream<IN,OUT> |
setConnectionId(String connectionId) |
ZeroMQChannelStream<IN,OUT> |
setSocket(org.zeromq.ZMQ.Socket socket) |
void |
subscribe(org.reactivestreams.Subscriber<? super IN> subscriber) |
String |
toString() |
getCapacity, getDecoder, getDispatcher, getEncoder, getEnvironment, writeBufferWith, writeWith
adaptiveConsume, adaptiveConsumeOn, after, batchConsume, batchConsumeOn, broadcast, broadcastOn, broadcastTo, buffer, buffer, buffer, buffer, buffer, buffer, buffer, buffer, buffer, buffer, buffer, cache, cancelSubscription, capacity, cast, combine, concatMap, concatWith, consume, consume, consume, consume, consume, consumeLater, consumeOn, consumeOn, consumeOn, count, count, decode, defaultIfEmpty, dematerialize, dispatchOn, dispatchOn, dispatchOn, distinct, distinct, distinctUntilChanged, distinctUntilChanged, downstreamSubscription, elapsed, elementAt, elementAtOrDefault, encode, env, exists, fanIn, filter, filter, finallyDo, flatMap, getTimer, groupBy, ignoreError, ignoreError, isReactivePull, join, joinWith, keepAlive, last, lift, log, log, map, materialize, merge, mergeWith, nest, next, observe, observeCancel, observeComplete, observeError, observeStart, observeSubscribe, onErrorResumeNext, onErrorResumeNext, onErrorReturn, onErrorReturn, onOverflowBuffer, onOverflowBuffer, onOverflowDrop, partition, partition, process, recover, reduce, reduce, repeat, repeat, repeatWhen, requestWhen, retry, retry, retry, retry, retryWhen, sample, sample, sample, sample, sample, sample, sampleFirst, sampleFirst, sampleFirst, sampleFirst, sampleFirst, sampleFirst, scan, scan, skip, skip, skip, skipWhile, skipWhile, sort, sort, sort, sort, split, split, startWith, startWith, startWith, subscribe, subscribeOn, subscribeOn, subscribeOn, switchMap, take, take, take, takeWhile, tap, throttle, throttle, timeout, timeout, timeout, timeout, timestamp, toBlockingQueue, toBlockingQueue, toList, toList, unbounded, when, window, window, window, window, window, window, window, window, window, window, window, zip, zipWith, zipWith
public ZeroMQChannelStream(Environment env, long prefetch, Dispatcher eventsDispatcher, InetSocketAddress remoteAddress, Codec<Buffer,IN,OUT> codec)
protected void doSubscribeWriter(org.reactivestreams.Publisher<? extends OUT> writer, org.reactivestreams.Subscriber<? super Void> postWriter)
doSubscribeWriter
in class ChannelStream<IN,OUT>
public void doDecoded(IN in)
doDecoded
in class ChannelStream<IN,OUT>
public void subscribe(org.reactivestreams.Subscriber<? super IN> subscriber)
public ZeroMQChannelStream<IN,OUT> setConnectionId(String connectionId)
public ZeroMQChannelStream<IN,OUT> setSocket(org.zeromq.ZMQ.Socket socket)
public InetSocketAddress remoteAddress()
ReactorChannel
public void close()
public ReactorChannel.ConsumerSpec on()
ReactorChannel
public org.zeromq.ZMQ.Socket delegate()
delegate
in class ChannelStream<IN,OUT>
Copyright © 2017. All rights reserved.