JAVA | Guava EventBus с использованием модели публикации/подписки

Java

Каталог статей серии


[TOC]


предисловие

EventBus — это механизм обработки событий Guava и реализация шаблона наблюдателя (модель производства/потребления).

Режим наблюдателя широко используется в нашей повседневной разработке.Например, в системе заказов изменения в статусе заказа или информации о логистике будут отправлять пользователям push-уведомления, SMS, уведомления продавцам, покупателям и т. д. в системе утверждения, процесс согласования заказов. Передача уведомит пользователя, инициировавшего согласование, руководителя согласования и т.д.

Режим Observer также поддерживается в JDK. Он уже существует в версии 1.0 Observer, но с быстрым обновлением версии Java его использование не изменилось. Многие библиотеки предоставляют более простые реализации, такие как Guava EventBus, RxJava, EventBus и т.д.

1. Зачем использовать шаблон Observer и преимущества EventBus?

Преимущества EventBus

  • По сравнению с Observer программирование простое и удобное
  • Синхронные и асинхронные операции и обработка исключений могут быть достигнуты с помощью пользовательских параметров.
  • Однопроцессное использование, без влияния на сеть

недостаток

  • Может использоваться только одним процессом
  • Аномальный перезапуск или выход из проекта не гарантирует сохранения сообщения.

Если вам нужно распределенное использование или вам нужно использоватьMQ

2. Шаги по использованию EventBus

1. Импортируйте библиотеку

Gradle

compile group: 'com.google.guava', name: 'guava', version: '29.0-jre'

Maven

<dependency>
    <groupId>com.google.guava</groupId>
    <artifactId>guava</artifactId>
    <version>29.0-jre</version>
</dependency>

После введения зависимостей здесь мы в основном используемcom.google.common.eventbus.EventBusкласс для работы, который обеспечиваетregister,unregister,postрегистрироваться, подписываться, отписываться и публиковать сообщения

public void register(Object object);

public void unregister(Object object);

public void post(Object event);

2. Синхронное использование

1. Сначала создайте EventBus

EventBus eventBus = new EventBus();

2. Создайте подписчика

В Guava EventBus подписка основана на типе параметра.Каждый метод подписки может иметь только один параметр, и он должен использовать@Subscribeлоготип

class EventListener {

  /**
   * 监听 Integer 类型的消息
   */
  @Subscribe
  public void listenInteger(Integer param) {
    System.out.println("EventListener#listenInteger ->" + param);
  }

  /**
   * 监听 String 类型的消息
   */
  @Subscribe
  public void listenString(String param) {
    System.out.println("EventListener#listenString ->" + param);
  }
}

3. Зарегистрируйтесь в EventBus и публикуйте сообщения

EventBus eventBus = new EventBus();

eventBus.register(new EventListener());

eventBus.post(1);
eventBus.post(2);
eventBus.post("3");

Результат бега есть

EventListener#listenInteger ->1
EventListener#listenInteger ->2
EventListener#listenString ->3

Мы можем создать несколько подписчиков для заполнения информации о подписке по мере необходимости, и если есть несколько подписчиков одного типа, будут выполнены все методы подписки.

Почему вы говорите, что это синхронно?

Событие Guava фактически использует пул потоков для обработки сообщений подписки.Из исходного кода видно, что когда мы используем конструктор по умолчанию для созданияEventBusкогдаexecutorзаMoreExecutors.directExecutor(), Конкретная реализация прямого вызоваRunnable#runметод, так что он все еще выполняется в том же потоке, поэтому операция по умолчанию все еще синхронна, и этот метод обработки также применим, так что он может не только отделить, но и разрешить выполнение метода в том же потоке. , например, обработка транзакций

Часть исходного кода EventBus

public class EventBus {
  private static final Logger logger = Logger.getLogger(EventBus.class.getName());
  private final String identifier;
  private final Executor executor;
  private final SubscriberExceptionHandler exceptionHandler;
  private final SubscriberRegistry subscribers;
  private final Dispatcher dispatcher;

  public EventBus() {
    this("default");
  }

  public EventBus(String identifier) {
    this(identifier, MoreExecutors.directExecutor(), Dispatcher.perThreadDispatchQueue(), EventBus.LoggingHandler.INSTANCE);
  }

  public EventBus(SubscriberExceptionHandler exceptionHandler) {
    this("default", MoreExecutors.directExecutor(), Dispatcher.perThreadDispatchQueue(), exceptionHandler);
  }

  EventBus(String identifier, Executor executor, Dispatcher dispatcher, SubscriberExceptionHandler exceptionHandler) {
    this.subscribers = new SubscriberRegistry(this);
    this.identifier = (String)Preconditions.checkNotNull(identifier);
    this.executor = (Executor)Preconditions.checkNotNull(executor);
    this.dispatcher = (Dispatcher)Preconditions.checkNotNull(dispatcher);
    this.exceptionHandler = (SubscriberExceptionHandler)Preconditions.checkNotNull(exceptionHandler);
  }
}

Часть исходного кода DirectExecutor

enum DirectExecutor implements Executor {
  INSTANCE;

  private DirectExecutor() {
  }

  public void execute(Runnable command) {
    command.run();
  }

  public String toString() {
    return "MoreExecutors.directExecutor()";
  }
}

3. Асинхронное использование

Из приведенного выше исходного кода видно, что пока исполнитель в методе построения заменен пулом потоков, Guava EventBus предоставляет упрощенное решение для упрощения операции.AsyncEventBus

EventBus eventBus = new AsyncEventBus(Executors.newCachedThreadPool());

Это позволяет асинхронно использовать

Исходный код AsyncEventBus

public class AsyncEventBus extends EventBus {
  public AsyncEventBus(String identifier, Executor executor) {
    super(identifier, executor, Dispatcher.legacyAsync(), LoggingHandler.INSTANCE);
  }

  public AsyncEventBus(Executor executor, SubscriberExceptionHandler subscriberExceptionHandler) {
    super("default", executor, Dispatcher.legacyAsync(), subscriberExceptionHandler);
  }

  public AsyncEventBus(Executor executor) {
    super("default", executor, Dispatcher.legacyAsync(), LoggingHandler.INSTANCE);
  }
}

4. Обработка исключений

Что делать, если во время обработки возникает исключение?EventBus все ещеAsyncEventBusНастраиваемыйSubscriberExceptionHandlerОбработчик вызывается при возникновении исключения, и я могу получить его из параметраexceptionПолучить информацию об исключении изcontextПолучить информацию о сообщении для конкретной обработки

Его интерфейс объявлен как

public interface SubscriberExceptionHandler {
  /** Handles exceptions thrown by subscribers. */
  void handleException(Throwable exception, SubscriberExceptionContext context);
}

Суммировать

На основе вышеизложенного мы можем определить некоторые типы сообщений для мониторинга и обработки различных сообщений.SubscriberExceptionHandlerЧтобы справиться с нештатными ситуациями, будь то синхронные или асинхронные, это можно сделать с легкостью.

Ссылаться на


白色兔子公众号图片