View Javadoc
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 jakarta.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   }