中文字幕av专区_日韩电影在线播放_精品国产精品久久一区免费式_av在线免费观看网站

溫馨提示×

溫馨提示×

您好,登錄后才能下訂單哦!

密碼登錄×
登錄注冊×
其他方式登錄
點擊 登錄注冊 即表示同意《億速云用戶服務條款》

spring與disruptor集成的示例分析

發布時間:2021-08-25 13:44:36 來源:億速云 閱讀:166 作者:小新 欄目:編程語言

這篇文章給大家分享的是有關spring與disruptor集成的示例分析的內容。小編覺得挺實用的,因此分享給大家做個參考,一起跟隨小編過來看看吧。

disruptor不過多介紹了,描述下當前的業務場景,兩個應用A,B,應用 A 向應用 B 傳遞數據 . 數據傳送比較快,如果用http直接push數據然后入庫,效率不高.有可能導致A應用比較大的壓力. 使用mq 太重量級,所以選擇了disruptor. 也可以使用Reactor

BaseQueueHelper.java

/**
 * lmax.disruptor 高效隊列處理模板. 支持初始隊列,即在init()前進行發布。
 *
 * 調用init()時才真正啟動線程開始處理 系統退出自動清理資源.
 *
 * @author xielongwang
 * @create 2018-01-18 下午3:49
 * @email xielong.wang@nvr-china.com
 * @description
 */
public abstract class BaseQueueHelper<D, E extends ValueWrapper<D>, H extends WorkHandler<E>> {

  /**
   * 記錄所有的隊列,系統退出時統一清理資源
   */
  private static List<BaseQueueHelper> queueHelperList = new ArrayList<BaseQueueHelper>();
  /**
   * Disruptor 對象
   */
  private Disruptor<E> disruptor;
  /**
   * RingBuffer
   */
  private RingBuffer<E> ringBuffer;
  /**
   * initQueue
   */
  private List<D> initQueue = new ArrayList<D>();

  /**
   * 隊列大小
   *
   * @return 隊列長度,必須是2的冪
   */
  protected abstract int getQueueSize();

  /**
   * 事件工廠
   *
   * @return EventFactory
   */
  protected abstract EventFactory<E> eventFactory();

  /**
   * 事件消費者
   *
   * @return WorkHandler[]
   */
  protected abstract WorkHandler[] getHandler();

  /**
   * 初始化
   */
  public void init() {
    ThreadFactory namedThreadFactory = new ThreadFactoryBuilder().setNameFormat("DisruptorThreadPool").build();
    disruptor = new Disruptor<E>(eventFactory(), getQueueSize(), namedThreadFactory, ProducerType.SINGLE, getStrategy());
    disruptor.setDefaultExceptionHandler(new MyHandlerException());
    disruptor.handleEventsWithWorkerPool(getHandler());
    ringBuffer = disruptor.start();

    //初始化數據發布
    for (D data : initQueue) {
      ringBuffer.publishEvent(new EventTranslatorOneArg<E, D>() {
        @Override
        public void translateTo(E event, long sequence, D data) {
          event.setValue(data);
        }
      }, data);
    }

    //加入資源清理鉤子
    synchronized (queueHelperList) {
      if (queueHelperList.isEmpty()) {
        Runtime.getRuntime().addShutdownHook(new Thread() {
          @Override
          public void run() {
            for (BaseQueueHelper baseQueueHelper : queueHelperList) {
              baseQueueHelper.shutdown();
            }
          }
        });
      }
      queueHelperList.add(this);
    }
  }

  /**
   * 如果要改變線程執行優先級,override此策略. YieldingWaitStrategy會提高響應并在閑時占用70%以上CPU,
   * 慎用SleepingWaitStrategy會降低響應更減少CPU占用,用于日志等場景.
   *
   * @return WaitStrategy
   */
  protected abstract WaitStrategy getStrategy();

  /**
   * 插入隊列消息,支持在對象init前插入隊列,則在隊列建立時立即發布到隊列處理.
   */
  public synchronized void publishEvent(D data) {
    if (ringBuffer == null) {
      initQueue.add(data);
      return;
    }
    ringBuffer.publishEvent(new EventTranslatorOneArg<E, D>() {
      @Override
      public void translateTo(E event, long sequence, D data) {
        event.setValue(data);
      }
    }, data);
  }

  /**
   * 關閉隊列
   */
  public void shutdown() {
    disruptor.shutdown();
  }
}

EventFactory.java

/**
 * @author xielongwang
 * @create 2018-01-18 下午6:24
 * @email xielong.wang@nvr-china.com
 * @description
 */
public class EventFactory implements com.lmax.disruptor.EventFactory<SeriesDataEvent> {

  @Override
  public SeriesDataEvent newInstance() {
    return new SeriesDataEvent();
  }
}

MyHandlerException.java

public class MyHandlerException implements ExceptionHandler {

  private Logger logger = LoggerFactory.getLogger(MyHandlerException.class);

  /*
   * (non-Javadoc) 運行過程中發生時的異常
   *
   * @see
   * com.lmax.disruptor.ExceptionHandler#handleEventException(java.lang.Throwable
   * , long, java.lang.Object)
   */
  @Override
  public void handleEventException(Throwable ex, long sequence, Object event) {
    ex.printStackTrace();
    logger.error("process data error sequence ==[{}] event==[{}] ,ex ==[{}]", sequence, event.toString(), ex.getMessage());
  }

  /*
   * (non-Javadoc) 啟動時的異常
   *
   * @see
   * com.lmax.disruptor.ExceptionHandler#handleOnStartException(java.lang.
   * Throwable)
   */
  @Override
  public void handleOnStartException(Throwable ex) {
    logger.error("start disruptor error ==[{}]!", ex.getMessage());
  }

  /*
   * (non-Javadoc) 關閉時的異常
   *
   * @see
   * com.lmax.disruptor.ExceptionHandler#handleOnShutdownException(java.lang
   * .Throwable)
   */
  @Override
  public void handleOnShutdownException(Throwable ex) {
    logger.error("shutdown disruptor error ==[{}]!", ex.getMessage());
  }
}

SeriesData.java (代表應用A發送給應用B的消息)

public class SeriesData {
  private String deviceInfoStr;
  public SeriesData() {
  }

  public SeriesData(String deviceInfoStr) {
    this.deviceInfoStr = deviceInfoStr;
  }

  public String getDeviceInfoStr() {
    return deviceInfoStr;
  }

  public void setDeviceInfoStr(String deviceInfoStr) {
    this.deviceInfoStr = deviceInfoStr;
  }

  @Override
  public String toString() {
    return "SeriesData{" +
        "deviceInfoStr='" + deviceInfoStr + '\'' +
        '}';
  }
}

SeriesDataEvent.java

public class SeriesDataEvent extends ValueWrapper<SeriesData> {
}

SeriesDataEventHandler.java

public class SeriesDataEventHandler implements WorkHandler<SeriesDataEvent> {
  private Logger logger = LoggerFactory.getLogger(SeriesDataEventHandler.class);
  @Autowired
  private DeviceInfoService deviceInfoService;

  @Override
  public void onEvent(SeriesDataEvent event) {
    if (event.getValue() == null || StringUtils.isEmpty(event.getValue().getDeviceInfoStr())) {
      logger.warn("receiver series data is empty!");
    }
    //業務處理
    deviceInfoService.processData(event.getValue().getDeviceInfoStr());
  }
}

SeriesDataEventQueueHelper.java

@Component
public class SeriesDataEventQueueHelper extends BaseQueueHelper<SeriesData, SeriesDataEvent, SeriesDataEventHandler> implements InitializingBean {
  private static final int QUEUE_SIZE = 1024;
  @Autowired
  private List<SeriesDataEventHandler> seriesDataEventHandler;

  @Override
  protected int getQueueSize() {
    return QUEUE_SIZE;
  }

  @Override
  protected com.lmax.disruptor.EventFactory eventFactory() {
    return new EventFactory();
  }

  @Override
  protected WorkHandler[] getHandler() {
    int size = seriesDataEventHandler.size();
    SeriesDataEventHandler[] paramEventHandlers = (SeriesDataEventHandler[]) seriesDataEventHandler.toArray(new SeriesDataEventHandler[size]);
    return paramEventHandlers;
  }

  @Override
  protected WaitStrategy getStrategy() {
    return new BlockingWaitStrategy();
    //return new YieldingWaitStrategy();
  }

  @Override
  public void afterPropertiesSet() throws Exception {
    this.init();
  }
}

ValueWrapper.java

public abstract class ValueWrapper<T> {
  private T value;
  public ValueWrapper() {}
  public ValueWrapper(T value) {
    this.value = value;
  }

  public T getValue() {
    return value;
  }

  public void setValue(T value) {
    this.value = value;
  }
}

DisruptorConfig.java

@Configuration
@ComponentScan(value = {"com.portal.disruptor"})
//多實例幾個消費者
public class DisruptorConfig {

  /**
   * smsParamEventHandler1
   *
   * @return SeriesDataEventHandler
   */
  @Bean
  public SeriesDataEventHandler smsParamEventHandler1() {
    return new SeriesDataEventHandler();
  }

  /**
   * smsParamEventHandler2
   *
   * @return SeriesDataEventHandler
   */
  @Bean
  public SeriesDataEventHandler smsParamEventHandler2() {
    return new SeriesDataEventHandler();
  }

  /**
   * smsParamEventHandler3
   *
   * @return SeriesDataEventHandler
   */
  @Bean
  public SeriesDataEventHandler smsParamEventHandler3() {
    return new SeriesDataEventHandler();
  }


  /**
   * smsParamEventHandler4
   *
   * @return SeriesDataEventHandler
   */
  @Bean
  public SeriesDataEventHandler smsParamEventHandler4() {
    return new SeriesDataEventHandler();
  }

  /**
   * smsParamEventHandler5
   *
   * @return SeriesDataEventHandler
   */
  @Bean
  public SeriesDataEventHandler smsParamEventHandler5() {
    return new SeriesDataEventHandler();
  }
}

測試

  //注入SeriesDataEventQueueHelper消息生產者
  @Autowired
  private SeriesDataEventQueueHelper seriesDataEventQueueHelper;

  @RequestMapping(value = "/data", method = RequestMethod.POST, produces = MediaType.APPLICATION_JSON_VALUE)
  public DataResponseVo<String> receiverDeviceData(@RequestBody String deviceData) {
    long startTime1 = System.currentTimeMillis();

    if (StringUtils.isEmpty(deviceData)) {
      logger.info("receiver data is empty !");
      return new DataResponseVo<String>(400, "failed");
    }
    seriesDataEventQueueHelper.publishEvent(new SeriesData(deviceData));
    long startTime2 = System.currentTimeMillis();
    logger.info("receiver data ==[{}] millisecond ==[{}]", deviceData, startTime2 - startTime1);
    return new DataResponseVo<String>(200, "success");
  }

應用A通過/data 接口把數據發送到應用B ,然后通過seriesDataEventQueueHelper 把消息發給disruptor隊列,消費者去消費,整個過程對不會堵塞應用A. 可接受消息丟失, 可以通過擴展SeriesDataEventQueueHelper來達到對disruptor隊列的監控

感謝各位的閱讀!關于“spring與disruptor集成的示例分析”這篇文章就分享到這里了,希望以上內容可以對大家有一定的幫助,讓大家可以學到更多知識,如果覺得文章不錯,可以把它分享出去讓更多的人看到吧!

向AI問一下細節

免責聲明:本站發布的內容(圖片、視頻和文字)以原創、轉載和分享為主,文章觀點不代表本網站立場,如果涉及侵權請聯系站長郵箱:is@yisu.com進行舉報,并提供相關證據,一經查實,將立刻刪除涉嫌侵權內容。

AI

平武县| 密云县| 个旧市| 阿巴嘎旗| 竹山县| 清流县| 河津市| 永清县| 安陆市| 花莲市| 铅山县| 兴仁县| 洪泽县| 洪江市| 彰化市| 嘉荫县| 宜阳县| 河池市| 南华县| 赤水市| 长宁县| 广汉市| 揭西县| 始兴县| 邵武市| 巴彦县| 昌黎县| 香河县| 沾益县| 玉山县| 宁城县| 泾源县| 枣阳市| 阿图什市| 申扎县| 界首市| 介休市| 永兴县| 宣化县| 巴东县| 洛浦县|