From 101f088e22c9e1c6d88847e1f7ceff4e0ab6ef1c Mon Sep 17 00:00:00 2001 From: Nikita Koksharov Date: Mon, 5 Nov 2018 15:58:53 +0300 Subject: [PATCH] RedissonBlockingQueueRx added --- .../redisson/rx/RedissonBlockingQueueRx.java | 84 +++++++++++++++++++ 1 file changed, 84 insertions(+) create mode 100644 redisson/src/main/java/org/redisson/rx/RedissonBlockingQueueRx.java diff --git a/redisson/src/main/java/org/redisson/rx/RedissonBlockingQueueRx.java b/redisson/src/main/java/org/redisson/rx/RedissonBlockingQueueRx.java new file mode 100644 index 000000000..0e3a8416a --- /dev/null +++ b/redisson/src/main/java/org/redisson/rx/RedissonBlockingQueueRx.java @@ -0,0 +1,84 @@ +/** + * Copyright 2018 Nikita Koksharov + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.redisson.rx; + +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; + +import org.redisson.RedissonBlockingQueue; +import org.redisson.api.RFuture; + +import io.netty.util.concurrent.Future; +import io.netty.util.concurrent.FutureListener; +import io.reactivex.Flowable; +import io.reactivex.functions.Action; +import io.reactivex.functions.LongConsumer; +import io.reactivex.processors.ReplayProcessor; + +/** + * + * @author Nikita Koksharov + * + * @param - value type + */ +public class RedissonBlockingQueueRx extends RedissonListRx { + + private final RedissonBlockingQueue queue; + + public RedissonBlockingQueueRx(RedissonBlockingQueue queue) { + super(queue); + this.queue = queue; + } + + public Flowable takeElements() { + final ReplayProcessor p = ReplayProcessor.create(); + return p.doOnRequest(new LongConsumer() { + @Override + public void accept(long n) throws Exception { + final AtomicLong counter = new AtomicLong(n); + final AtomicReference> futureRef = new AtomicReference>(); + take(p, counter, futureRef); + p.doOnCancel(new Action() { + @Override + public void run() throws Exception { + futureRef.get().cancel(true); + } + }); + } + }); + } + + private void take(final ReplayProcessor p, final AtomicLong counter, final AtomicReference> futureRef) { + RFuture future = queue.takeAsync(); + futureRef.set(future); + future.addListener(new FutureListener() { + @Override + public void operationComplete(Future future) throws Exception { + if (!future.isSuccess()) { + p.onError(future.cause()); + return; + } + + p.onNext(future.getNow()); + if (counter.decrementAndGet() == 0) { + p.onComplete(); + } + + take(p, counter, futureRef); + } + }); + } +}