org.reactivestreams.Processor类的使用及代码示例

x33g5p2x  于2022-01-26 转载在 其他  
字(3.3k)|赞(0)|评价(0)|浏览(333)

本文整理了Java中org.reactivestreams.Processor类的一些代码示例,展示了Processor类的具体用法。这些代码示例主要来源于Github/Stackoverflow/Maven等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Processor类的具体详情如下:
包路径:org.reactivestreams.Processor
类名称: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();
}

相关文章