ICode9

精准搜索请尝试: 精确搜索
首页 > 编程语言> 文章详细

reactor-core – java.lang.IllegalStateException:队列已满?!在热发布者(ConnectableFlux)上

2019-07-27 15:01:51  阅读:159  来源: 互联网

标签:project-reactor java spring rx-java reactive-programming


到目前为止我一直在使用RxJava,但我开始使用来自projectreactor.io的reactor-core,因为它遵循反应流规范.

在下面的测试中,我创建了一个生成随机数的热Flux(ConnectableFlux).我立即连接()它预取256个值(我可以在日志中看到它们实际上有258个).我等待5秒钟来模拟订阅者直到一段时间后才会订阅.

主线程唤醒后,RnApp订阅了ConnectableFlux,randomNumberGenerator.subscribe(new RnApp());.然后调用RnApp.onSubscribe()并请求10个元素.之后,java.lang.IllegalStateException:Queue full?!异常被引发(调用RnApp.onError()),为什么?

订户:

public class RnApp implements Subscriber<Float>{

    private Subscription subscription;
    private List<Float> randomNumbers = new ArrayList<Float>();

    @Override
    public void onComplete() {
        System.out.println("onComplete");
    }

    @Override
    public void one rror(Throwable err) {
        err.printStackTrace();
    }

    @Override
    public void onNext(Float f) {
        if(this.randomNumbers.size()>=10){
            this.subscription.cancel();
        }else{
            this.randomNumbers.add(f);
        }
    }

    @Override
    public void onSubscribe(Subscription subs) {
        this.subscription = subs;
        this.subscription.request(10);
    }
}

出版商测试:

@Test
public void randomNumberReading() throws InterruptedException {

    CountDownLatch latch = new CountDownLatch(1);
    ConnectableFlux<Float> randomNumberGenerator = ConnectableFlux.<Float>create( (c) -> {
        SecureRandom sr = new SecureRandom();
        int i = 1;
        while(true){
            try {
                Thread.sleep(10);
            } catch (Exception e) {
                e.printStackTrace();
            }
            System.out.println("-----------------------------------------------------"+(i++));
            c.onNext(sr.nextFloat());
        }
    }).log().subscribeOn(Computations.concurrent()).publish();

    randomNumberGenerator.connect();

    Thread.sleep(5000);

    randomNumberGenerator.subscribe(new RnApp());

    latch.await();
}

日志:

11:12:05.125 [main] DEBUG r.core.util.Logger$LoggerFactory - Using Slf4j logging framework
11:12:05.363 [concurrent-1] INFO  reactor.core.publisher.FluxLog -  onSubscribe(io.pivotal.literx.Part10SubscribeOnPublishOn$$Lambda$1/1586600255@29d4caeb)
11:12:05.371 [concurrent-1] INFO  reactor.core.publisher.FluxLog -  request(256)
-----------------------------------------------------1
11:12:06.000 [concurrent-1] INFO  reactor.core.publisher.FluxLog -  onNext(0.39189225)
-----------------------------------------------------2
...
-----------------------------------------------------257
11:12:08.683 [concurrent-1] INFO  reactor.core.publisher.FluxLog -  onNext(0.34729618)
-----------------------------------------------------258
11:12:08.697 [concurrent-1] INFO  reactor.core.publisher.FluxLog -  onNext(0.7729547)
java.lang.IllegalStateException: Queue full?!
    at reactor.core.publisher.FluxPublish$State.onNext(FluxPublish.java:246)
    at reactor.core.publisher.FluxSubscribeOn$SubscribeOnPipeline.onNext(FluxSubscribeOn.java:134)
    at reactor.core.publisher.FluxLog$LoggerBarrier.doNext(FluxLog.java:130)
    at reactor.core.subscriber.SubscriberBarrier.onNext(SubscriberBarrier.java:85)
    at reactor.core.subscriber.SubscriberWithContext.onNext(SubscriberWithContext.java:92)
    at io.pivotal.literx.Part10SubscribeOnPublishOn.lambda$1(Part10SubscribeOnPublishOn.java:132)
    at reactor.core.publisher.FluxGenerate$ForEachBiConsumer.accept(FluxGenerate.java:145)
    at reactor.core.publisher.FluxGenerate$ForEachBiConsumer.accept(FluxGenerate.java:114)
    at reactor.core.publisher.FluxGenerate$SubscriberProxy.request(FluxGenerate.java:245)
    at reactor.core.subscriber.SubscriberBarrier.doRequest(SubscriberBarrier.java:146)
    at reactor.core.publisher.FluxLog$LoggerBarrier.doRequest(FluxLog.java:160)
    at reactor.core.subscriber.SubscriberBarrier.request(SubscriberBarrier.java:135)
    at reactor.core.util.DeferredSubscription.set(DeferredSubscription.java:71)
    at reactor.core.publisher.FluxSubscribeOn$SubscribeOnPipeline.onSubscribe(FluxSubscribeOn.java:129)
    at reactor.core.publisher.FluxLog$LoggerBarrier.doOnSubscribe(FluxLog.java:122)
    at reactor.core.subscriber.SubscriberBarrier.onSubscribe(SubscriberBarrier.java:67)
    at reactor.core.publisher.FluxGenerate.subscribe(FluxGenerate.java:72)
    at reactor.core.publisher.FluxLog.subscribe(FluxLog.java:67)
    at reactor.core.publisher.FluxSubscribeOn$SourceSubscribeTask.run(FluxSubscribeOn.java:363)
    at reactor.core.publisher.Computations$ProcessorWorker.onNext(Computations.java:919)
    at reactor.core.publisher.Computations$ProcessorWorker.onNext(Computations.java:883)
    at reactor.core.publisher.WorkQueueProcessor$QueueSubscriberLoop.run(WorkQueueProcessor.java:842)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source)
    at java.lang.Thread.run(Unknown Source)

解决方法:

与RxJava一样,如果您使用的是create(),那么您可以自己处理取消和背压.您可以使用标准运算符构建生成器:

ConnectableFlux<Double> secureRandomFlux = Flux.using(
    () -> new SecureRandom(),
    sr -> Flux.interval(10, TimeUnit.MILLISECONDS)
          .map(v -> sr.nextDouble())
          .onBackpressureDrop()
    sr -> { }
).publish();

标签:project-reactor,java,spring,rx-java,reactive-programming
来源: https://codeday.me/bug/20190727/1555146.html

本站声明: 1. iCode9 技术分享网(下文简称本站)提供的所有内容,仅供技术学习、探讨和分享;
2. 关于本站的所有留言、评论、转载及引用,纯属内容发起人的个人观点,与本站观点和立场无关;
3. 关于本站的所有言论和文字,纯属内容发起人的个人观点,与本站观点和立场无关;
4. 本站文章均是网友提供,不完全保证技术分享内容的完整性、准确性、时效性、风险性和版权归属;如您发现该文章侵犯了您的权益,可联系我们第一时间进行删除;
5. 本站为非盈利性的个人网站,所有内容不会用来进行牟利,也不会利用任何形式的广告来间接获益,纯粹是为了广大技术爱好者提供技术内容和技术思想的分享性交流网站。

专注分享技术,共同学习,共同进步。侵权联系[81616952@qq.com]

Copyright (C)ICode9.com, All Rights Reserved.

ICode9版权所有