本文整理了Java中reactor.core.publisher.Operators.unboundedOrPrefetch()
方法的一些代码示例,展示了Operators.unboundedOrPrefetch()
的具体用法。这些代码示例主要来源于Github
/Stackoverflow
/Maven
等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Operators.unboundedOrPrefetch()
方法的具体详情如下:
包路径:reactor.core.publisher.Operators
类名称:Operators
方法名:unboundedOrPrefetch
暂无
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.setOnce(S, this, s)) {
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.setOnce(S, this, s)) {
this.cachedSubscription = s;
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.setOnce(S, this, s)) {
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.setOnce(S, this, s)) {
s.request(Operators.unboundedOrPrefetch(batchSize));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.validate(this.s, s)) {
this.s = s;
actual.onSubscribe(this);
if (cancelled == 0) {
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.validate(this.s, s)) {
this.s = s;
actual.onSubscribe(this);
s.request(Operators.unboundedOrPrefetch(maxConcurrency));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.validate(this.s, s)) {
this.s = s;
actual.onSubscribe(this);
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.setOnce(S, this, s)) {
if (s instanceof Fuseable.QueueSubscription) {
@SuppressWarnings("unchecked") Fuseable.QueueSubscription<T> f =
(Fuseable.QueueSubscription<T>) s;
int m = f.requestFusion(Fuseable.ANY | Fuseable.THREAD_BARRIER);
if (m == Fuseable.SYNC) {
sourceMode = m;
queue = f;
drain();
return;
}
if (m == Fuseable.ASYNC) {
sourceMode = m;
queue = f;
s.request(Operators.unboundedOrPrefetch(prefetch));
return;
}
}
queue = parent.queueSupplier.get();
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.setOnce(S, this, s)) {
if (s instanceof QueueSubscription) {
@SuppressWarnings("unchecked") QueueSubscription<T> qs =
(QueueSubscription<T>) s;
int m = qs.requestFusion(Fuseable.ANY);
if (m == Fuseable.SYNC) {
sourceMode = m;
queue = qs;
done = true;
connected = true;
drain();
return;
}
if (m == Fuseable.ASYNC) {
sourceMode = m;
queue = qs;
connected = true;
s.request(Operators.unboundedOrPrefetch(prefetch));
return;
}
}
queue = queueSupplier.get();
connected = true;
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(final Subscription s) {
if (Operators.setOnce(S, this, s)) {
if (s instanceof Fuseable.QueueSubscription) {
@SuppressWarnings("unchecked") Fuseable.QueueSubscription<T> f =
(Fuseable.QueueSubscription<T>) s;
int m = f.requestFusion(Fuseable.ANY);
if (m == Fuseable.SYNC) {
sourceMode = m;
queue = f;
drain();
return;
}
else if (m == Fuseable.ASYNC) {
sourceMode = m;
queue = f;
s.request(Operators.unboundedOrPrefetch(prefetch));
return;
}
}
queue = Queues.<T>get(prefetch).get();
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.setOnce(SUBSCRIPTION, this, s)) {
if (s instanceof QueueSubscription) {
@SuppressWarnings("unchecked")
QueueSubscription<R> qs = (QueueSubscription<R>) s;
int m = qs.requestFusion(Fuseable.ANY | Fuseable.THREAD_BARRIER);
if (m == Fuseable.SYNC) {
fusionMode = m;
queue = qs;
done = true;
parent.innerComplete(this);
return;
}
if (m == Fuseable.ASYNC) {
fusionMode = m;
queue = qs;
s.request(Operators.unboundedOrPrefetch(prefetch));
return;
}
}
queue = Queues.<R>get(prefetch).get();
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@SuppressWarnings("unchecked")
@Override
public void onSubscribe(Subscription s) {
if (Operators.setOnce(S, this, s)) {
if (s instanceof Fuseable.QueueSubscription) {
Fuseable.QueueSubscription<T> f = (Fuseable.QueueSubscription<T>) s;
int m = f.requestFusion(Fuseable.ANY | Fuseable.THREAD_BARRIER);
if (m == SYNC) {
sourceMode = SYNC;
queue = f;
done = true;
parent.drain();
return;
}
else if (m == ASYNC) {
sourceMode = ASYNC;
queue = f;
}
else {
queue = queueSupplier.get();
}
}
else {
queue = queueSupplier.get();
}
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.validate(this.s, s)) {
this.s = s;
if (s instanceof QueueSubscription) {
@SuppressWarnings("unchecked") QueueSubscription<T> f =
(QueueSubscription<T>) s;
int m = f.requestFusion(Fuseable.ANY | Fuseable.THREAD_BARRIER);
if (m == Fuseable.SYNC) {
sourceMode = Fuseable.SYNC;
queue = f;
done = true;
actual.onSubscribe(this);
return;
}
if (m == Fuseable.ASYNC) {
sourceMode = Fuseable.ASYNC;
queue = f;
actual.onSubscribe(this);
s.request(Operators.unboundedOrPrefetch(prefetch));
return;
}
}
queue = queueSupplier.get();
actual.onSubscribe(this);
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.validate(this.s, s)) {
this.s = s;
if (s instanceof QueueSubscription) {
@SuppressWarnings("unchecked") QueueSubscription<T> f =
(QueueSubscription<T>) s;
int m = f.requestFusion(Fuseable.ANY | Fuseable.THREAD_BARRIER);
if (m == Fuseable.SYNC) {
sourceMode = Fuseable.SYNC;
queue = f;
done = true;
actual.onSubscribe(this);
return;
}
if (m == Fuseable.ASYNC) {
sourceMode = Fuseable.ASYNC;
queue = f;
actual.onSubscribe(this);
s.request(Operators.unboundedOrPrefetch(prefetch));
return;
}
}
queue = queueSupplier.get();
actual.onSubscribe(this);
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.validate(this.s, s)) {
this.s = s;
if (s instanceof QueueSubscription) {
@SuppressWarnings("unchecked") QueueSubscription<T> qs =
(QueueSubscription<T>) s;
int m = qs.requestFusion(Fuseable.ANY);
if (m == Fuseable.SYNC) {
fusionMode = m;
this.queue = qs;
done = true;
actual.onSubscribe(this);
return;
}
else if (m == Fuseable.ASYNC) {
fusionMode = m;
this.queue = qs;
actual.onSubscribe(this);
s.request(Operators.unboundedOrPrefetch(prefetch));
return;
}
}
queue = queueSupplier.get();
actual.onSubscribe(this);
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
s.request(Operators.unboundedOrPrefetch(prefetch));
s.request(Operators.unboundedOrPrefetch(prefetch));
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
if (Operators.setOnce(S, this, s)) {
if (s instanceof Fuseable.QueueSubscription) {
@SuppressWarnings("unchecked") Fuseable.QueueSubscription<R> f =
(Fuseable.QueueSubscription<R>) s;
int m = f.requestFusion(Fuseable.ANY | Fuseable.THREAD_BARRIER);
if (m == Fuseable.SYNC) {
sourceMode = Fuseable.SYNC;
queue = f;
done = true;
parent.drain();
return;
}
if (m == Fuseable.ASYNC) {
sourceMode = Fuseable.ASYNC;
queue = f;
}
// NONE is just fall-through as the queue will be created on demand
}
s.request(Operators.unboundedOrPrefetch(prefetch));
}
}
代码示例来源:origin: reactor/reactor-core
@Test
public void unboundedOrPrefetch() {
assertThat(Operators.unboundedOrPrefetch(10))
.as("bounded")
.isEqualTo(10L);
assertThat(Operators.unboundedOrPrefetch(Integer.MAX_VALUE))
.as("unbounded")
.isEqualTo(Long.MAX_VALUE);
}
代码示例来源:origin: reactor/reactor-core
s.request(Operators.unboundedOrPrefetch(prefetch));
代码示例来源:origin: reactor/reactor-core
@Override
public void onSubscribe(Subscription s) {
Subscription a = this.s;
if (a == Operators.cancelledSubscription()) {
s.cancel();
}
if (a != null) {
s.cancel();
Operators.reportSubscriptionSet();
return;
}
if (S.compareAndSet(this, null, s)) {
s.request(Operators.unboundedOrPrefetch(prefetch));
return;
}
a = this.s;
if (a != Operators.cancelledSubscription()) {
s.cancel();
Operators.reportSubscriptionSet();
}
}
内容来源于网络,如有侵权,请联系作者删除!