|
5 | 5 | import io.reactivex.Single; |
6 | 6 | import io.reactivex.schedulers.Schedulers; |
7 | 7 | import io.reflectoring.reactive.batch.MessageHandler.Result; |
8 | | -import org.reactivestreams.Subscriber; |
9 | | -import org.reactivestreams.Subscription; |
10 | | - |
11 | | -import java.util.ArrayList; |
12 | | -import java.util.List; |
| 8 | +import java.util.concurrent.Executors; |
13 | 9 | import java.util.concurrent.LinkedBlockingDeque; |
14 | 10 | import java.util.concurrent.ThreadPoolExecutor; |
15 | 11 | import java.util.concurrent.TimeUnit; |
| 12 | +import org.reactivestreams.Subscriber; |
| 13 | +import org.reactivestreams.Subscription; |
16 | 14 |
|
17 | 15 | public class ReactiveBatchProcessor { |
18 | 16 |
|
19 | | - private final static Logger logger = new Logger(); |
20 | | - |
21 | | - private final int MESSAGE_BATCHES = 10; |
22 | | - |
23 | | - private final int BATCH_SIZE = 3; |
24 | | - |
25 | | - private final int THREADS = 4; |
26 | | - |
27 | | - private final MessageHandler messageHandler = new MessageHandler(); |
28 | | - |
29 | | - public void start() { |
30 | | - |
31 | | - Scheduler threadPoolScheduler = threadPoolScheduler(THREADS, 10); |
32 | | - |
33 | | - retrieveMessageBatches() |
34 | | - .doOnNext(batch -> logger.log(batch.toString())) |
35 | | - .flatMap(batch -> Flowable.fromIterable(batch.getMessages())) |
36 | | - .flatMapSingle(message -> Single.defer(() -> Single.just(messageHandler.handleMessage(message))) |
37 | | - .doOnSuccess(result -> logger.log("message handled")) |
38 | | - .subscribeOn(threadPoolScheduler)) |
39 | | - .subscribeWith(subscriber()); |
40 | | - } |
41 | | - |
42 | | - private Subscriber<MessageHandler.Result> subscriber() { |
43 | | - return new Subscriber<>() { |
44 | | - private Subscription subscription; |
45 | | - |
46 | | - @Override |
47 | | - public void onSubscribe(Subscription subscription) { |
48 | | - this.subscription = subscription; |
49 | | - subscription.request(THREADS); |
50 | | - logger.log("subscribed"); |
51 | | - } |
52 | | - |
53 | | - @Override |
54 | | - public void onNext(Result message) { |
55 | | - subscription.request(THREADS); |
56 | | - } |
57 | | - |
58 | | - @Override |
59 | | - public void onError(Throwable t) { |
60 | | - logger.log("error"); |
61 | | - } |
62 | | - |
63 | | - @Override |
64 | | - public void onComplete() { |
65 | | - logger.log("completed"); |
66 | | - } |
67 | | - }; |
68 | | - } |
69 | | - |
70 | | - private Scheduler threadPoolScheduler(int poolSize, int queueSize) { |
71 | | - return Schedulers.from(new ThreadPoolExecutor( |
72 | | - poolSize, |
73 | | - poolSize, |
74 | | - 0L, |
75 | | - TimeUnit.SECONDS, |
76 | | - new LinkedBlockingDeque<>(queueSize), |
77 | | - new RetryRejectedExecutionHandler() |
78 | | - )); |
79 | | - } |
80 | | - |
81 | | - public boolean allMessagesProcessed() { |
82 | | - return this.messageHandler.getProcessedMessages() == MESSAGE_BATCHES * BATCH_SIZE; |
83 | | - } |
84 | | - |
85 | | - private Flowable<MessageBatch> retrieveMessageBatches() { |
86 | | - return Flowable.range(1, MESSAGE_BATCHES) |
87 | | - .map(this::messageBatch); |
88 | | - } |
89 | | - |
90 | | - private MessageBatch messageBatch(int batchNumber) { |
91 | | - List<Message> messages = new ArrayList<>(); |
92 | | - for (int i = 1; i <= BATCH_SIZE; i++) { |
93 | | - messages.add(new Message(String.format("%d-%d", batchNumber, i))); |
94 | | - } |
95 | | - return new MessageBatch(messages); |
96 | | - } |
97 | | - |
| 17 | + private final static Logger logger = new Logger(); |
| 18 | + |
| 19 | + private final int threads; |
| 20 | + |
| 21 | + private final MessageHandler messageHandler; |
| 22 | + |
| 23 | + private final MessageSource messageSource; |
| 24 | + |
| 25 | + public ReactiveBatchProcessor( |
| 26 | + MessageSource messageSource, |
| 27 | + MessageHandler messageHandler, |
| 28 | + int threads) { |
| 29 | + this.messageSource = messageSource; |
| 30 | + this.threads = threads; |
| 31 | + this.messageHandler = messageHandler; |
| 32 | + } |
| 33 | + |
| 34 | + public void start() { |
| 35 | + |
| 36 | + Scheduler scheduler = threadPoolScheduler(threads, 10); |
| 37 | + |
| 38 | + messageSource.getMessageBatches() |
| 39 | + .subscribeOn(Schedulers.from(Executors.newSingleThreadExecutor())) |
| 40 | + .doOnNext(batch -> logger.log(batch.toString())) |
| 41 | + .flatMap(batch -> Flowable.fromIterable(batch.getMessages())) |
| 42 | + .flatMapSingle(m -> Single.defer(() -> Single.just(m) |
| 43 | + .map(messageHandler::handleMessage)) |
| 44 | + .subscribeOn(scheduler)) |
| 45 | + .subscribeWith(subscriber()); |
| 46 | + } |
| 47 | + |
| 48 | + private Subscriber<MessageHandler.Result> subscriber() { |
| 49 | + return new Subscriber<>() { |
| 50 | + private Subscription subscription; |
| 51 | + |
| 52 | + @Override |
| 53 | + public void onSubscribe(Subscription subscription) { |
| 54 | + this.subscription = subscription; |
| 55 | + subscription.request(threads); |
| 56 | + logger.log("subscribed"); |
| 57 | + } |
| 58 | + |
| 59 | + @Override |
| 60 | + public void onNext(Result message) { |
| 61 | + subscription.request(threads); |
| 62 | + } |
| 63 | + |
| 64 | + @Override |
| 65 | + public void onError(Throwable t) { |
| 66 | + logger.log("error"); |
| 67 | + } |
| 68 | + |
| 69 | + @Override |
| 70 | + public void onComplete() { |
| 71 | + logger.log("completed"); |
| 72 | + } |
| 73 | + }; |
| 74 | + } |
| 75 | + |
| 76 | + private Scheduler threadPoolScheduler(int poolSize, int queueSize) { |
| 77 | + return Schedulers.from(new ThreadPoolExecutor( |
| 78 | + poolSize, |
| 79 | + poolSize, |
| 80 | + 0L, |
| 81 | + TimeUnit.SECONDS, |
| 82 | + new LinkedBlockingDeque<>(queueSize) |
| 83 | + )); |
| 84 | + } |
98 | 85 |
|
99 | 86 | } |
0 commit comments