io.reactivex.subjects.SingleSubject#create ( )源码实例Demo

下面列出了io.reactivex.subjects.SingleSubject#create ( ) 实例代码,或者点击链接到github查看源代码,也可以在右侧发表评论。

源代码1 项目: crnk-framework   文件: AppServer.java
public void start() {
	VertxOptions options = new VertxOptions();
	options.setMaxEventLoopExecuteTime(Long.MAX_VALUE);

	SingleSubject waitSubject = SingleSubject.create();
	Handler<AsyncResult<String>> completionHandler = event -> {
		if (event.succeeded()) {
			waitSubject.onSuccess(event.result());
		} else {
			event.cause().printStackTrace();
			System.exit(0);
		}
	};


	vertx = Vertx.vertx(options);
	vertx.deployVerticle(vehicle, completionHandler);
	waitSubject.blockingGet();
}
 
源代码2 项目: marvel   文件: SearchInteractorImpl.java
@Override
public Single<CharactersResponse> loadCharacter(String query,
                                                String privateKey,
                                                String publicKey,
                                                long timestamp) {
    if (characterSubscription == null || characterSubscription.isDisposed()) {
        characterSubject = SingleSubject.create();

        // generate hash using timestamp and API keys
        String hash = HashGenerator.generate(timestamp, privateKey, publicKey);

        characterSubscription = api.getCharacters(query, publicKey, hash, timestamp)
                .subscribeOn(scheduler.backgroundThread())
                .subscribe(characterSubject::onSuccess);
    }

    return characterSubject.hide();
}
 
源代码3 项目: RxCentralBle   文件: CorePeripheral.java
@Override
public Single<byte[]> read(UUID svc, UUID chr) {
  final SingleSubject<Pair<UUID, byte[]>> scopedSubject = SingleSubject.create();

  return scopedSubject
      .filter(chrPair -> chrPair.first.equals(chr))
      .flatMapSingle(chrPair -> Single.just(chrPair.second))
      .doOnSubscribe(disposable -> processRead(scopedSubject, svc, chr))
      .doFinally(() -> clearReadSubject(scopedSubject));
}
 
源代码4 项目: RxCentralBle   文件: CorePeripheral.java
@Override
public Completable write(UUID svc, UUID chr, byte[] data) {
  final SingleSubject<UUID> scopedSubject = SingleSubject.create();

  return scopedSubject
      .filter(writeChr -> writeChr.equals(chr))
      .ignoreElement()
      .doOnSubscribe(disposable -> processWrite(scopedSubject, svc, chr, data))
      .doFinally(() -> clearWriteSubject(scopedSubject));
}
 
源代码5 项目: RxCentralBle   文件: CorePeripheral.java
@Override
public Completable registerNotification(UUID svc, UUID chr, @Nullable Preprocessor preprocessor) {
  final SingleSubject<UUID> scopedSubject = SingleSubject.create();

  return scopedSubject
      .filter(regChr -> regChr.equals(chr))
      .ignoreElement()
      .doOnSubscribe(
          disposable -> processRegisterNotification(scopedSubject, svc, chr, preprocessor))
      .doFinally(() -> clearRegisterNotificationSubject(scopedSubject));
}
 
源代码6 项目: RxCentralBle   文件: CorePeripheral.java
@Override
public Completable unregisterNotification(UUID svc, UUID chr) {
  final SingleSubject<UUID> scopedSubject = SingleSubject.create();

  return scopedSubject
      .filter(regChr -> regChr.equals(chr))
      .ignoreElement()
      .doOnSubscribe(disposable -> processUnregisterNotification(scopedSubject, svc, chr))
      .doOnComplete(() -> preprocessorMap.remove(chr))
      .doFinally(() -> clearRegisterNotificationSubject(scopedSubject));
}
 
源代码7 项目: RxCentralBle   文件: CorePeripheral.java
@TargetApi(21)
@Override
public Single<Integer> requestMtu(int mtu) {
  if (Build.VERSION.SDK_INT < Build.VERSION_CODES.LOLLIPOP) {
    return Single.error(new PeripheralError(PeripheralError.Code.MINIMUM_SDK_UNSUPPORTED));
  }

  final SingleSubject<Integer> scopedSubject = SingleSubject.create();

  return scopedSubject
      .doOnSubscribe(disposable -> processRequestMtu(scopedSubject, mtu))
      .doFinally(() -> clearRequestMtuSubject(scopedSubject));
}
 
源代码8 项目: RxCentralBle   文件: CorePeripheral.java
@Override
public Single<Integer> readRssi() {
  SingleSubject<Integer> scopedSubject = SingleSubject.create();

  return scopedSubject
      .doOnSubscribe(disposable -> processReadRssi(scopedSubject))
      .doFinally(() -> clearReadRssiSubject(scopedSubject));
}
 
源代码9 项目: crnk-framework   文件: CrnkVertxHandler.java
private Mono getBody(VertxRequestContext vertxRequestContext) {
    HttpServerRequest serverRequest = vertxRequestContext.getServerRequest();
    SingleSubject<VertxRequestContext> bodySubject = SingleSubject.create();
    Handler<Buffer> bodyHandler = (event) -> {
        vertxRequestContext.setRequestBody(event.toString().getBytes(StandardCharsets.UTF_8));
        bodySubject.onSuccess(vertxRequestContext);
    };
    serverRequest.bodyHandler(bodyHandler);
    return Mono.from(bodySubject.toFlowable());
}
 
源代码10 项目: crnk-framework   文件: VertxTestContainer.java
@Override
public void stop() {
    SingleSubject waitSubject = SingleSubject.create();
    Handler<AsyncResult<Void>> completionHandler = event -> waitSubject.onSuccess("test");
    vertx.close(completionHandler);
    waitSubject.blockingGet();

    vehicle.testModule.clear();
    vertx = null;
    vehicle = null;
    port = -1;
}
 
源代码11 项目: akarnokd-misc   文件: ConcurrentSingleCache.java
public Single<Entity> getEntity(String guid) {
    SingleSubject<Entity> e = map.get(guid);
    if (e == null) {
        e = SingleSubject.create();
        SingleSubject<Entity> f = map.putIfAbsent(guid, e);
        if (f == null) {
            longLoadFromDatabase(guid).subscribe(e);
        } else {
            e = f;
        }
    }
    return e;
}
 
源代码12 项目: crnk-framework   文件: AppServer.java
public void stop() {
	SingleSubject waitSubject = SingleSubject.create();
	Handler<AsyncResult<Void>> completionHandler = event -> waitSubject.onSuccess("test");
	vertx.close(completionHandler);
	waitSubject.blockingGet();
}
 
源代码13 项目: akarnokd-misc   文件: SingleObserveOnRaceTest.java
@Test
public void race() throws Exception {
    Worker w = Schedulers.newThread().createWorker();
    try {
        for (int i = 0; i < 1000; i++) {
            Integer[] value = { 0, 0 };
            
            TestObserver<Integer> to = new TestObserver<Integer>() {
                @Override
                public void onSuccess(Integer v) {
                    value[1] = value[0];
                    super.onSuccess(v);
                }
            };
            
            SingleSubject<Integer> subj = SingleSubject.create();
            
            subj.observeOn(Schedulers.single())
            .onTerminateDetach()
            .subscribe(to);
            
            AtomicInteger wip = new AtomicInteger(2);
            CountDownLatch cdl = new CountDownLatch(2);
            
            w.schedule(() -> {
                if (wip.decrementAndGet() != 0) {
                    while (wip.get() != 0);
                }
                subj.onSuccess(1);
                cdl.countDown();
            });
            
            Schedulers.single().scheduleDirect(() -> {
                if (wip.decrementAndGet() != 0) {
                    while (wip.get() != 0);
                }
                to.cancel();
                value[0] = null;
                cdl.countDown();
            });

            cdl.await();
            
            Assert.assertNotNull(value[1]);
        }
    } finally {
        w.dispose();
    }
}