- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
Rx 世界很新,我需要实现以下行为:
一旦我有至少 N 个项目,或者如果自 以来经过了 T 时间,我需要 observable 来收集值并将它们作为列表发出。组的第一项被发射了。
我一遍又一遍地阅读文档,很确定它会使用
buffer(timespan, unit, count[, scheduler])
最佳答案
这个问题的缓冲方面可以通过多播技巧来实现,但我发现为它编写一个运算符要容易得多,因此数据和上下文位于同一个可访问的位置:
public final class OperatorBufferFirst<T> implements Operator<List<T>, T> {
final Scheduler scheduler;
final long timeout;
final TimeUnit unit;
final int maxSize;
public OperatorBufferFirst(
long timeout, TimeUnit unit,
Scheduler scheduler, int maxSize) {
this.timeout = timeout;
this.unit = unit;
this.scheduler = scheduler;
this.maxSize = maxSize;
}
@Override
public Subscriber<? super T> call(
Subscriber<? super List<T>> t) {
BufferSubscriber<T> parent = new BufferSubscriber<>(
new SerializedSubscriber<>(t),
timeout, unit,
scheduler.createWorker(), maxSize);
t.add(parent);
return parent;
}
static final class BufferSubscriber<T>
extends Subscriber<T> {
final Subscriber<? super List<T>> actual;
final Scheduler.Worker w;
final long timeout;
final TimeUnit unit;
final int maxSize;
final SerialSubscription timer;
List<T> buffer;
long index;
public BufferSubscriber(
Subscriber<? super List<T>> actual,
long timeout,
TimeUnit unit,
Scheduler.Worker w,
int maxSize) {
this.actual = actual;
this.timeout = timeout;
this.unit = unit;
this.w = w;
this.maxSize = maxSize;
this.timer = new SerialSubscription();
this.buffer = new ArrayList<>();
this.add(timer);
this.add(w);
}
@Override
public void onNext(T t) {
List<T> b;
boolean startTimer = false;
boolean emit = false;
long idx;
synchronized (this) {
b = buffer;
b.add(t);
idx = index;
int n = b.size();
if (n == 1) {
startTimer = true;
} else
if (n < maxSize) {
return;
} else {
buffer = new ArrayList<>();
index = ++idx;
emit = true;
}
}
if (startTimer) {
final long fidx = idx;
timer.set(w.schedule(() -> timeout(fidx), timeout, unit));
}
if (emit) {
timer.set(Subscriptions.unsubscribed());
actual.onNext(b);
}
}
@Override
public void onError(Throwable e) {
actual.onError(e);
}
@Override
public void onCompleted() {
timer.unsubscribe();
List<T> b;
synchronized (this) {
b = buffer;
buffer = null;
index++;
}
if (!b.isEmpty()) {
actual.onNext(b);
}
actual.onCompleted();
}
public void timeout(long idx) {
List<T> b;
synchronized (this) {
b = buffer;
if (idx != index) {
return;
}
buffer = new ArrayList<>();
index = idx + 1;
}
actual.onNext(b);
}
}
public static void main(String[] args) {
TestScheduler s = Schedulers.test();
PublishSubject<Integer> source = PublishSubject.create();
source.lift(new OperatorBufferFirst<>(1, TimeUnit.SECONDS, s, 3))
.subscribe(System.out::println, Throwable::printStackTrace,
() -> System.out.println("Done"));
source.onNext(1);
source.onNext(2);
source.onNext(3);
source.onNext(4);
s.advanceTimeBy(1, TimeUnit.SECONDS);
source.onNext(5);
source.onNext(6);
s.advanceTimeBy(1, TimeUnit.SECONDS);
s.advanceTimeBy(1, TimeUnit.SECONDS);
source.onNext(7);
source.onCompleted();
}
}
synchronized (this) {
b = buffer;
idx = index;
if (t != FLUSH) {
b.add(t);
int n = b.size();
if (n == 1) {
startTimer = true;
} else
if (n < maxSize) {
return;
} else {
buffer = new ArrayList<>();
index = ++idx;
emit = true;
}
} else {
buffer = new ArrayList<>();
index = ++idx;
emit = true;
}
}
关于system.reactive - 自第一个新组元素以来超时的 Rx 缓冲区,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/33775889/
我想知道有没有可能做 new PrintWriter(new BufferedWriter(new PrintWriter(s.getOutputStream, true))) 在 Java 中,s
我正在尝试使用 ConcurrentHashMap 初始化 ConcurrentHashMap private final ConcurrentHashMap > myMulitiConcurrent
我只是想知道两个不同的新对象初始化器之间是否有任何区别,还是仅仅是语法糖。 因此: Dim _StreamReader as New Streamreader(mystream) 与以下内容不同: D
在 C++ 中,以下两种动态对象创建之间的确切区别是什么: A* pA = new A; A* pA = new A(); 我做了一些测试,但似乎在这两种情况下,都调用了默认构造函数,并且只调用了它。
我已经阅读了其他帖子,但它们没有解决我的问题。环境为VB 2008(2.0 Framework)下面的代码在 xslt.Load 行导致 XSLT 编译错误下面是错误的输出。我将 XSLT 作为字符串
我想知道为什么alert(new Boolean(false))打印 false 而不是打印对象,因为 new Boolean 应该返回对象。如果我使用 console.log(new Boolean
本文实例讲述了Python装饰器用法。分享给大家供大家参考,具体如下: 写装饰器 装饰器只不过是一种函数,接收被装饰的可调用对象作为它的唯一参数,然后返回一个可调用对象(就像前面的简单例子) 注
我可以编写 YAML header 来使用 knit 为 R Markdown 文件生成多种输出格式吗?我无法重现 the original question with this title 的答案中
我可以编写一个YAML标头以使用knitr为R Markdown文件生成多种输出格式吗?我无法重现the original question with this title答案中描述的功能。 这个降价
我正在使用vars package可视化脉冲响应。示例: library(vars) Canada % names ir % `$`(irf) %>% `[[`(variables[e])) %>%
我有一个容器类,它有一个通用参数,该参数被限制到某个基类。提供给泛型的类型是基类约束的子类。子类使用方法隐藏(新)来更改基类方法的行为(不,我不能将其设为虚拟,因为它不是我的代码)。我的问题是"new
Java 在提示! cannot find symbol symbol : constructor Bar() location: class Bar JPanel panel =
在我的应用程序中,一个新的 Activity 从触摸按钮(而不是点击)开始,而且我没有抬起手指并希望在新的 Activity 中跟踪触摸的 Action 。第二个 Activity 中的触摸监听器不响
已关闭。此问题旨在寻求有关书籍、工具、软件库等的建议。不符合Stack Overflow guidelines .它目前不接受答案。 我们不允许提问寻求书籍、工具、软件库等的推荐。您可以编辑问题,
和我的last question ,我的程序无法检测到一个短语并将其与第一行以外的任何行匹配。但是,我已经解决并回答了。但现在我需要一个新的 def函数,它删除某个(给定 refName )联系人及其
这个问题在这里已经有了答案: Horizontal list items (7 个答案) 关闭 9 年前。
我想创建一个新的 float 类型,大小为 128 位,指数为 4 字节(32 位),小数为 12 字节(96 位),我该怎么做输入 C++,我将能够在其中进行输入、输出、+、-、*、/操作。 [我正
我在放置引用计数指针的实例时遇到问题 类到我的数组类中。使用调试器,似乎永远不会调用构造函数(这会扰乱引用计数并导致行中出现段错误)! 我的 push_back 函数是: void push_back
我在我们的代码库中发现了经典的新建/删除不匹配错误,如下所示: char *foo = new char[10]; // do something delete foo; // instead of
A *a = new A(); 这是创建一个指针还是一个对象? 我是一个 c++ 初学者,所以我想了解这个区别。 最佳答案 两者:您创建了一个新的 A 实例(一个对象),并创建了一个指向它的名为 a
我是一名优秀的程序员,十分优秀!