- iOS/Objective-C 元类和类别
- objective-c - -1001 错误,当 NSURLSession 通过 httpproxy 和/etc/hosts
- java - 使用网络类获取 url 地址
- ios - 推送通知中不播放声音
到目前为止,我一直在使用 RxJava,但我开始使用来自 projectreactor.io 的 reactor-core,因为它符合 react 流规范。
在下面的测试中,我创建了一个生成随机数的热通量 (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 onError(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();
关于java - react 堆核心 - java.lang.IllegalStateException : Queue full? !在 Hot Publisher (ConnectableFlux) 上,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/37000396/
我试图弄清楚以下模块正在做什么。 import Queue import multiprocessing import threading class BufferedReadQueue(Queue.
如果我使用 Queue.Queue,那么我的 read() 函数不起作用,为什么?但是,如果我使用 multiprocessing.Queue,它运行良好: from multiprocessing
我正在寻找比我在文档中找到的更多关于 Python 队列实现的见解。 根据我的理解,如果我在这方面有误,请原谅我的无知: queue.Queue():通过内存中的基本数组实现,因此不能在多个进程之间共
当我使用多处理模块(Windows 上的 Python 2.7)中的队列代替 Queue.Queue 时,我的程序没有完全关闭。 最终,我想使用 multiprocessing.Process 处理
阅读了大量的 JavaScript 事件循环教程,我看到了不同的术语来标识队列存储消息,当调用堆栈为空时,事件循环准备好获取消息: 队列 消息队列 事件队列 我找不到规范的术语来识别它。 甚至 MDN
我收到错误消息“类型队列不接受参数”。当我将更改队列行替换为 PriorityQueue 时,此错误消失并且编译正常。有什么区别以及如何将其更改为编译队列和常规队列? import java.util
如何将项目返回到 queue.Queue?如果任务失败,这在线程或多处理中很有用,这样任务就不会丢失。 docs for queue.Queue.get()说函数可以“从队列中删除并返回一个项目”,但
如何在多个 queue.Queue 上进行“选择”同时? Golang 有 desired feature及其 channel : select { case i1 = 声明。 线程:queue 模
http://docs.python.org/2/library/queue.html#Queue.Queue.put 这似乎是一个幼稚的问题,但我在文档和谷歌搜索中都没有找到答案,那么这些方法是线程
这可能是个愚蠢的问题,但我对与 .dequeue() 和 $.queue() 一起使用的 .queue() 感到困惑> 或 jquery.queue()。 它们是否相同,如果是,为什么 jquery
我正在尝试创建一个线程化的 tcp 流处理程序类线程和主线程对话,但是 Queue.Queue 也没有做我需要的,服务器从另一个程序接收数据,我只想传递它进入主线程进行处理这里是我到目前为止的代码:
The principal challenge of multi-threaded applications is coordinating threads that share data or ot
在Queue模块的queue类中,有几个方法,分别是qsize、empty 和 full,其文档声称它们“不可靠”。 他们到底有什么不可靠的地方? 我确实注意到 on the Python docs网
我需要一个队列,多个线程可以将内容放入其中,并且多个线程可以从中读取。 Python 至少有两个队列类,Queue.Queue 和 collections.deque,前者似乎在内部使用后者。两者都在
明天我将介绍我选择进程内消息队列实现的基本原理,但我无法阐明我的推理。我的合作设计者提议我们实现一个简单的异步队列,只使用基本的作业列表和互斥锁来控制访问,我建议在嵌入式模式下使用 ActiveMQ。
在 scala 中定义了一个特征: trait Queue[T] Queue 是一种类型吗?或其他东西,例如类型构造函数? 来自 http://artima.com/pins1ed/type-para
我看到 SML/NJ 包含一个队列结构。我不知道如何使用它。如何使用 SML/NJ 提供的附加库? 最佳答案 Queue structure SML '97 未指定,但它存在于 SML/NJ 的顶级环
我是 D3 和 JavaScript 的新手。 我试图理解其中的 queue.js。 我已经完成了 this关联。但是仍然无法清楚地了解 queue.await() 和 queue.awaitAll(
所以我试图在我的 main.cpp 文件中调用一个函数,但我得到“错误:没有匹配函数来调用‘Queue::Queue()。” 队列.h #ifndef QUEUE_H #define QUEUE_H
假设我有一个 10 行的二维 numpy 数组 例如 array([[ 23425. , 521331.40625], [ 23465. , 521246.03125],
我是一名优秀的程序员,十分优秀!