rx.subjects.Subject.onBackpressureBuffer()方法的使用及代码示例

x33g5p2x  于2022-01-30 转载在 其他  
字(2.4k)|赞(0)|评价(0)|浏览(121)

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

Subject.onBackpressureBuffer介绍

暂无

代码示例

代码示例来源:origin: PipelineAI/pipeline

/* package */ HystrixRequestEventsStream() {
  writeOnlyRequestEventsSubject = PublishSubject.create();
  readOnlyRequestEvents = writeOnlyRequestEventsSubject.onBackpressureBuffer(1024);
}

代码示例来源:origin: PipelineAI/pipeline

/* package */ HystrixThreadEventStream(Thread thread) {
  this.threadId = thread.getId();
  this.threadName = thread.getName();
  writeOnlyCommandStartSubject = PublishSubject.create();
  writeOnlyCommandCompletionSubject = PublishSubject.create();
  writeOnlyCollapserSubject = PublishSubject.create();
  writeOnlyCommandStartSubject
      .onBackpressureBuffer()
      .doOnNext(writeCommandStartsToShardedStreams)
      .unsafeSubscribe(Subscribers.empty());
  writeOnlyCommandCompletionSubject
      .onBackpressureBuffer()
      .doOnNext(writeCommandCompletionsToShardedStreams)
      .unsafeSubscribe(Subscribers.empty());
  writeOnlyCollapserSubject
      .onBackpressureBuffer()
      .doOnNext(writeCollapserExecutionsToShardedStreams)
      .unsafeSubscribe(Subscribers.empty());
}

代码示例来源:origin: dswarm/dswarm

public Observable<Collection<Triple>> getObservable() {

    return tripleSubject.onBackpressureBuffer(10000).filter(NOT_NULL);
  }
}

代码示例来源:origin: com.netflix.hystrix/hystrix-core

/* package */ HystrixRequestEventsStream() {
  writeOnlyRequestEventsSubject = PublishSubject.create();
  readOnlyRequestEvents = writeOnlyRequestEventsSubject.onBackpressureBuffer(1024);
}

代码示例来源:origin: com.netflix.hystrix/hystrix-core

/* package */ HystrixThreadEventStream(Thread thread) {
  this.threadId = thread.getId();
  this.threadName = thread.getName();
  writeOnlyCommandStartSubject = PublishSubject.create();
  writeOnlyCommandCompletionSubject = PublishSubject.create();
  writeOnlyCollapserSubject = PublishSubject.create();
  writeOnlyCommandStartSubject
      .onBackpressureBuffer()
      .doOnNext(writeCommandStartsToShardedStreams)
      .unsafeSubscribe(Subscribers.empty());
  writeOnlyCommandCompletionSubject
      .onBackpressureBuffer()
      .doOnNext(writeCommandCompletionsToShardedStreams)
      .unsafeSubscribe(Subscribers.empty());
  writeOnlyCollapserSubject
      .onBackpressureBuffer()
      .doOnNext(writeCollapserExecutionsToShardedStreams)
      .unsafeSubscribe(Subscribers.empty());
}

相关文章