本文整理了Java中org.reactivestreams.Processor
类的一些代码示例,展示了Processor
类的具体用法。这些代码示例主要来源于Github
/Stackoverflow
/Maven
等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Processor
类的具体详情如下:
包路径:org.reactivestreams.Processor
类名称:Processor
[英]A Processor represents a processing stage—which is both a Subscriberand a Publisher and obeys the contracts of both.
[中]处理器代表一个处理阶段,它既是订阅者又是发布者,并遵守两者的合同。
代码示例来源:origin: ReactiveX/RxJava
/**
* Emit the given values and complete the Processor.
* @param <T> the value type
* @param p the target processor
* @param values the values to emit
*/
public static <T> void emit(Processor<T, ?> p, T... values) {
for (T v : values) {
p.onNext(v);
}
p.onComplete();
}
代码示例来源:origin: micronaut-projects/micronaut-core
@Override
public void onError(Throwable error) {
processor.onError(error);
}
代码示例来源:origin: micronaut-projects/micronaut-core
@Override
public void subscribe(Subscriber<? super WebSocketFrame> subscriber) {
processor.subscribe(subscriber);
}
代码示例来源:origin: micronaut-projects/micronaut-core
@Override
public void onComplete() {
processor.onComplete();
}
}
代码示例来源:origin: micronaut-projects/micronaut-core
@Override
public void onNext(WebSocketFrame webSocketFrame) {
processor.onNext(webSocketFrame);
}
代码示例来源:origin: io.apptik.rhub/roxy-rs
@Override
public void onError(Throwable t) {
if (tePolicy.equals(WRAP)) {
proc.onNext(new Event.ErrorEvent(t));
} else if (tePolicy.equals(PASS)) {
proc.onError(t);
}
}
代码示例来源:origin: micronaut-projects/micronaut-core
@Override
public void onSubscribe(Subscription subscription) {
processor.onSubscribe(subscription);
}
代码示例来源:origin: ReactiveX/RxJava
@Override
public void onComplete() {
Processor<T, T> w = window;
if (w != null) {
window = null;
w.onComplete();
}
downstream.onComplete();
}
代码示例来源:origin: reactor/reactor-core
void sendNext(int count) {
for (int i = 0; i < count; i++) {
// System.out.println("XXXX " + x);
processor.onNext((x++) + "\n");
}
}
}
代码示例来源:origin: apptik/RHub
@Override
public void onError(Throwable t) {
if (tePolicy.equals(WRAP)) {
proc.onNext(new Event.ErrorEvent(t));
} else if (tePolicy.equals(PASS)) {
proc.onError(t);
}
}
代码示例来源:origin: akarnokd/RxJava2Extensions
@Override
public void onSubscribe(Subscription s) {
source.onSubscribe(s);
}
代码示例来源:origin: ReactiveX/RxJava
w.onNext(t);
w.onComplete();
代码示例来源:origin: ReactiveX/RxJava
@Override
public void onComplete() {
Processor<T, T> w = window;
if (w != null) {
window = null;
w.onComplete();
}
downstream.onComplete();
}
代码示例来源:origin: akarnokd/RxJava2Extensions
@Override
public void onNext(T t) {
source.onNext(t);
}
代码示例来源:origin: ReactiveX/RxJava
@Override
public void onError(Throwable t) {
Processor<T, T> w = window;
if (w != null) {
window = null;
w.onError(t);
}
downstream.onError(t);
}
代码示例来源:origin: apptik/RHub
@Override
public void onError(Throwable t) {
if (tePolicy.equals(WRAP)) {
proc.onNext(new Event.ErrorEvent(t));
cnt.decrementAndGet();
cnt.compareAndSet(0, mat.settings().maxInputBufferSize());
} else if (tePolicy.equals(PASS)) {
proc.onError(t);
}
}
代码示例来源:origin: redisson/redisson
@Override
public void subscribe(Subscriber<? super E> s) {
try {
processor.subscribe(s);
} catch (Throwable t) {
s.onError(t);
}
}
};
代码示例来源:origin: com.github.akarnokd/rxjava2-extensions
@Override
public void onSubscribe(Subscription s) {
source.onSubscribe(s);
}
代码示例来源:origin: redisson/redisson
w.onNext(t);
w.onComplete();
代码示例来源:origin: redisson/redisson
@Override
public void onComplete() {
Processor<T, T> w = window;
if (w != null) {
window = null;
w.onComplete();
}
actual.onComplete();
}
内容来源于网络,如有侵权,请联系作者删除!