Coverage Report - com.ctrip.framework.apollo.biz.message.ReleaseMessageScanner
 
Classes in this File Line Coverage Branch Coverage Complexity
ReleaseMessageScanner
83%
47/56
72%
13/18
0
 
 1  
 package com.ctrip.framework.apollo.biz.message;
 2  
 
 3  
 import com.google.common.collect.Lists;
 4  
 
 5  
 import com.ctrip.framework.apollo.biz.entity.ReleaseMessage;
 6  
 import com.ctrip.framework.apollo.biz.repository.ReleaseMessageRepository;
 7  
 import com.ctrip.framework.apollo.core.utils.ApolloThreadFactory;
 8  
 import com.dianping.cat.Cat;
 9  
 import com.dianping.cat.message.Message;
 10  
 import com.dianping.cat.message.Transaction;
 11  
 
 12  
 import org.slf4j.Logger;
 13  
 import org.slf4j.LoggerFactory;
 14  
 import org.springframework.beans.factory.InitializingBean;
 15  
 import org.springframework.beans.factory.annotation.Autowired;
 16  
 import org.springframework.core.env.Environment;
 17  
 import org.springframework.util.CollectionUtils;
 18  
 
 19  
 import java.util.List;
 20  
 import java.util.Objects;
 21  
 import java.util.concurrent.Executors;
 22  
 import java.util.concurrent.ScheduledExecutorService;
 23  
 import java.util.concurrent.TimeUnit;
 24  
 
 25  
 /**
 26  
  * @author Jason Song(song_s@ctrip.com)
 27  
  */
 28  
 public class ReleaseMessageScanner implements InitializingBean {
 29  1
   private static final Logger logger = LoggerFactory.getLogger(ReleaseMessageScanner.class);
 30  
   private static final int DEFAULT_SCAN_INTERVAL_IN_MS = 1000;
 31  
   @Autowired
 32  
   private Environment env;
 33  
   @Autowired
 34  
   private ReleaseMessageRepository releaseMessageRepository;
 35  
   private int databaseScanInterval;
 36  
   private List<ReleaseMessageListener> listeners;
 37  
   private ScheduledExecutorService executorService;
 38  
   private long maxIdScanned;
 39  
 
 40  1
   public ReleaseMessageScanner() {
 41  1
     listeners = Lists.newLinkedList();
 42  2
     executorService = Executors.newScheduledThreadPool(1, ApolloThreadFactory
 43  1
         .create("ReleaseMessageScanner", true));
 44  1
   }
 45  
 
 46  
   @Override
 47  
   public void afterPropertiesSet() throws Exception {
 48  1
     populateDataBaseInterval();
 49  1
     maxIdScanned = loadLargestMessageId();
 50  2
     executorService.scheduleWithFixedDelay((Runnable) () -> {
 51  18
       Transaction transaction = Cat.newTransaction("Apollo.ReleaseMessageScanner", "scanMessage");
 52  
       try {
 53  18
         scanMessages();
 54  18
         transaction.setStatus(Message.SUCCESS);
 55  0
       } catch (Throwable ex) {
 56  0
         transaction.setStatus(ex);
 57  0
         logger.error("Scan and send message failed", ex);
 58  
       } finally {
 59  18
         transaction.complete();
 60  18
       }
 61  19
     }, getDatabaseScanIntervalMs(), getDatabaseScanIntervalMs(), TimeUnit.MILLISECONDS);
 62  
 
 63  1
   }
 64  
 
 65  
   /**
 66  
    * add message listeners for release message
 67  
    * @param listener
 68  
    */
 69  
   public void addMessageListener(ReleaseMessageListener listener) {
 70  2
     if (!listeners.contains(listener)) {
 71  2
       listeners.add(listener);
 72  
     }
 73  2
   }
 74  
 
 75  
   /**
 76  
    * Scan messages, continue scanning until there is no more messages
 77  
    */
 78  
   private void scanMessages() {
 79  18
     boolean hasMoreMessages = true;
 80  36
     while (hasMoreMessages && !Thread.currentThread().isInterrupted()) {
 81  18
       hasMoreMessages = scanAndSendMessages();
 82  
     }
 83  18
   }
 84  
 
 85  
   /**
 86  
    * scan messages and send
 87  
    *
 88  
    * @return whether there are more messages
 89  
    */
 90  
   private boolean scanAndSendMessages() {
 91  
     //current batch is 500
 92  18
     List<ReleaseMessage> releaseMessages =
 93  18
         releaseMessageRepository.findFirst500ByIdGreaterThanOrderByIdAsc(maxIdScanned);
 94  18
     if (CollectionUtils.isEmpty(releaseMessages)) {
 95  16
       return false;
 96  
     }
 97  2
     fireMessageScanned(releaseMessages);
 98  2
     int messageScanned = releaseMessages.size();
 99  2
     maxIdScanned = releaseMessages.get(messageScanned - 1).getId();
 100  2
     return messageScanned == 500;
 101  
   }
 102  
 
 103  
   /**
 104  
    * find largest message id as the current start point
 105  
    * @return current largest message id
 106  
    */
 107  
   private long loadLargestMessageId() {
 108  1
     ReleaseMessage releaseMessage = releaseMessageRepository.findTopByOrderByIdDesc();
 109  1
     return releaseMessage == null ? 0 : releaseMessage.getId();
 110  
   }
 111  
 
 112  
   /**
 113  
    * Notify listeners with messages loaded
 114  
    * @param messages
 115  
    */
 116  
   private void fireMessageScanned(List<ReleaseMessage> messages) {
 117  2
     for (ReleaseMessage message : messages) {
 118  2
       for (ReleaseMessageListener listener : listeners) {
 119  
         try {
 120  3
           listener.handleMessage(message, Topics.APOLLO_RELEASE_TOPIC);
 121  0
         } catch (Throwable ex) {
 122  0
           Cat.logError(ex);
 123  0
           logger.error("Failed to invoke message listener {}", listener.getClass(), ex);
 124  3
         }
 125  3
       }
 126  2
     }
 127  2
   }
 128  
 
 129  
   private void populateDataBaseInterval() {
 130  1
     databaseScanInterval = DEFAULT_SCAN_INTERVAL_IN_MS;
 131  
     try {
 132  1
       String interval = env.getProperty("apollo.message-scan.interval");
 133  1
       if (!Objects.isNull(interval)) {
 134  1
         databaseScanInterval = Integer.parseInt(interval);
 135  
       }
 136  0
     } catch (Throwable ex) {
 137  0
       Cat.logError(ex);
 138  0
       logger.error("Load apollo message scan interval from system property failed", ex);
 139  1
     }
 140  1
   }
 141  
 
 142  
   private int getDatabaseScanIntervalMs() {
 143  2
     return databaseScanInterval;
 144  
   }
 145  
 }