spring与disruptor集成的简单示例

 更新时间:2018年02月26日 09:51:19   作者:Muroidea  
本篇文章主要介绍了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 Boot加载properties和yml配置文件

    详解Spring Boot加载properties和yml配置文件

    本篇文章主要介绍了详解Spring Boot加载properties和yml配置文件,具有一定的参考价值,感兴趣的小伙伴们可以参考一下
    2017-04-04
  • mybatis-plus配置日志两种实现方式

    mybatis-plus配置日志两种实现方式

    这篇文章主要给大家介绍了关于mybatis-plus配置日志两种实现方式的相关资料,Mybatis-plus集成了日志框架,可以将程序运行时产生的日志进行记录,方便开发人员进行问题排查,需要的朋友可以参考下
    2023-09-09
  • 移动指定文件夹内的全部文件

    移动指定文件夹内的全部文件

    移动指定文件夹内的全部文件
    2009-01-01
  • 远程debug调试入门

    远程debug调试入门

    这篇文章主要介绍了Eclipse的Debug调试技巧大全(总结),小编觉得挺不错的,现在分享给大家,也给大家做个参考。一起跟随小编过来看看吧吗,希望能给你带来帮助
    2021-06-06
  • 策略模式:告别if else

    策略模式:告别if else

    你是不是还在写着大量的if else语句,if else 不仅难以维护不易扩展,而且使代码臃肿不堪,想不想让你的业务代码更加的健壮,更易扩展,那你一定要学一学今天的主角策略模式
    2021-06-06
  • java读取枚举类的值转成list和map方式

    java读取枚举类的值转成list和map方式

    这篇文章主要介绍了java读取枚举类的值转成list和map方式,具有很好的参考价值,希望对大家有所帮助。如有错误或未考虑完全的地方,望不吝赐教
    2023-07-07
  • 认识Java底层操作系统与并发基础

    认识Java底层操作系统与并发基础

    这篇文章主要介绍了认识Java底层操作系统与并发基础,文章围绕主题展开详细的内容介绍,具有一定的参考价值,需要的朋友可以参考一下
    2022-07-07
  • 如何从eureka获取服务的ip和端口号进行Http的调用

    如何从eureka获取服务的ip和端口号进行Http的调用

    这篇文章主要介绍了如何从eureka获取服务的ip和端口号进行Http的调用,具有很好的参考价值,希望对大家有所帮助。如有错误或未考虑完全的地方,望不吝赐教
    2022-03-03
  • springboot接收别人上传的本地视频实例代码

    springboot接收别人上传的本地视频实例代码

    本文通过实例代码给大家介绍了springboot接收别人上传的本地视频,代码简单易懂,非常不错,具有一定的参考借鉴价值,需要的朋友参考下吧
    2018-07-07
  • Java中防止XSS攻击的10种方法

    Java中防止XSS攻击的10种方法

    这篇文章主要给大家介绍了关于Java中防止XSS攻击的10种方法,XSS(Cross Site Scripting)攻击是一种常见的网络安全漏洞,文中给出了详细的代码示例,需要的朋友可以参考下
    2023-07-07

最新评论