1 /*
2 * *************************************************************************************************************************************************************
3 *
4 * TheseFoolishThings: Miscellaneous utilities
5 * http://tidalwave.it/projects/thesefoolishthings
6 *
7 * Copyright (C) 2009 - 2025 by Tidalwave s.a.s. (http://tidalwave.it)
8 *
9 * *************************************************************************************************************************************************************
10 *
11 * Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License.
12 * You may obtain a copy of the License at
13 *
14 * http://www.apache.org/licenses/LICENSE-2.0
15 *
16 * 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
17 * CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License.
18 *
19 * *************************************************************************************************************************************************************
20 *
21 * git clone https://bitbucket.org/tidalwave/thesefoolishthings-src
22 * git clone https://github.com/tidalwave-it/thesefoolishthings-src
23 *
24 * *************************************************************************************************************************************************************
25 */
26 package it.tidalwave.messagebus.spi;
27
28 import javax.annotation.Nonnull;
29 import it.tidalwave.messagebus.spi.MultiQueue.TopicAndMessage;
30 import lombok.Getter;
31 import lombok.Setter;
32 import lombok.ToString;
33 import lombok.extern.slf4j.Slf4j;
34
35 /***************************************************************************************************************************************************************
36 *
37 * An implementation of {@link MessageDelivery} that dispatches messages in a round-robin fashion, topic by topic.
38 * Each delivery is performed in a separated thread.
39 *
40 * @author Fabrizio Giudici
41 * @since 2.2
42 *
43 **************************************************************************************************************************************************************/
44 @Slf4j @ToString(of = "workers")
45 public class RoundRobinAsyncMessageDelivery implements MessageDelivery
46 {
47 @Nonnull
48 private SimpleMessageBus messageBusSupport;
49
50 @Getter @Setter
51 private int workers = 10;
52
53 private final MultiQueue multiQueue = new MultiQueue();
54
55 /***********************************************************************************************************************************************************
56 **********************************************************************************************************************************************************/
57 private final Runnable dispatcher = new Runnable()
58 {
59 @Override
60 public void run()
61 {
62 for (;;)
63 {
64 try
65 {
66 dispatchMessage(multiQueue.remove());
67 }
68 catch (InterruptedException e)
69 {
70 break;
71 }
72 }
73 }
74
75 private <T> void dispatchMessage (@Nonnull final TopicAndMessage<T> tam)
76 {
77 messageBusSupport.dispatchMessage(tam.getTopic(), tam.getMessage());
78 }
79 };
80
81 /***********************************************************************************************************************************************************
82 * {@inheritDoc}
83 **********************************************************************************************************************************************************/
84 @Override
85 public void initialize (@Nonnull final SimpleMessageBus messageBusSupport)
86 {
87 this.messageBusSupport = messageBusSupport;
88 final var executor = this.messageBusSupport.getExecutor();
89
90 for (var i = 0; i < workers; i++)
91 {
92 executor.execute(dispatcher);
93 }
94 }
95
96 /***********************************************************************************************************************************************************
97 * {@inheritDoc}
98 **********************************************************************************************************************************************************/
99 @Override
100 public <T> void deliverMessage (@Nonnull final Class<T> topic, @Nonnull final T message)
101 {
102 multiQueue.add(topic, message);
103 }
104 }