- 使用 Spring Initializr 创建 Spring Boot 应用程序
- 在Spring Boot中配置Cassandra
- 在 Spring Boot 上配置 Tomcat 连接池
- 将Camel消息路由到嵌入WildFly的Artemis上
消息中间件提供了系统之间的异步处理机制,比如在电商网站上支付订单之后,会触发库存计算,物流调度计算,甚至是营销人员绩效计算,报表统计等,诸如此类的操作一般会耗费比订单购买商品本身更多的时间,加之这样的操作没有即时的时效性要求,用户在下单之后完全没有必要等待电商后端做完所有操作才算成功,那么此时消息中间件是一种非常好的解决方案,用户下单支付之后即可向用户返回购买成功的通知,然后提交各种消息到消息中间件,这样注册在消息中间件的其他系统就可以顺利地接受订单通知了,然后执行各自的业务逻辑。消息中间件主要用于解决进程之间消息异步处理的解决方案。
Bus 接口:对外提供几种主要的使用方式,比如 post 方法用来发送 Event,register 方法用来注册 Evnet 接收者(Subscriber)接受相应事件,EventBus 采用同步的方式推送 Event,AsyncEventBus 采用异步的方式(Thread-Per-Message)推送 Event。
Register 注册表:主要用来记录对应 Subscriber 以及受理消息的回调方法,回调方法用注解 @Subscribe 来标识。
Dispatcher:主要用来将 event 广播给注册表中监听了 topic 的 Subscriber。
package concurrent.eventbus;
/**
* @className: Bus
* @description: 定义了 EventBus 的所有使用方法
* @date: 2022/5/13
* @author: cakin
*/
public interface Bus {
// 将某个对象注册到 Bus 上,从此之后该类就成为了 Subscriber 了
void register(Object subscriber);
// 将某个对象从 Bus 上取消注册,取消注册之后就不会再接受到来自 Bus 的任何消息
void unregister(Object subscriber);
// 提交 Event 到默认的 topic
void post(Object event);
// 提交 Event 到指定的 topic
void post(Object Event, String topic);
// 关闭该 bus
void close();
// 返回 Bus 的名称标识
String getBusName();
}
package concurrent.eventbus;
import java.util.concurrent.Executor;
/**
* @className: EventBus
* @description: 它实现了 Bus 所有的功能,采用的是同步的方式
* @date: 2022/5/13
* @author: cakin
*/
public class EventBus implements Bus {
// 用于维护 Subscriber 的注册表
private final Registry registry = new Registry();
// Event Bus 的名字
private String busName;
// 默认的 Event Bus 的名字
private final static String DEFAULT_BUS_NAME = "default";
// 默认的 Event Bus 的名字
private final static String DEFAULT_TOPIC = "default-topic";
// 用于分发广播消息到各个 Subscriber 的类
private final Dispatcher dispatcher;
public EventBus() {
this(DEFAULT_BUS_NAME, null, Dispatcher.SEQ_EXECUTOR_SERVICE);
}
public EventBus(String busName) {
this(busName, null, Dispatcher.SEQ_EXECUTOR_SERVICE);
}
public EventBus(String busName, EventExceptionHandler eventExceptionHandler, Executor executor) {
this.busName = busName;
this.dispatcher = Dispatcher.newDispatcher(eventExceptionHandler, executor);
}
public EventBus(EventExceptionHandler eventExceptionHandler) {
this(DEFAULT_BUS_NAME, eventExceptionHandler, Dispatcher.SEQ_EXECUTOR_SERVICE);
}
@Override
public void register(Object subscriber) {
this.registry.bind(subscriber);
}
@Override
public void unregister(Object subscriber) {
this.registry.unbind(subscriber);
}
@Override
public void post(Object event) {
this.post(event, DEFAULT_TOPIC);
}
@Override
public void post(Object event, String topic) {
this.dispatcher.dispatch(this, registry, event, topic);
}
@Override
public void close() {
this.dispatcher.close();
}
@Override
public String getBusName() {
return null;
}
}
package concurrent.eventbus;
import java.util.concurrent.ThreadPoolExecutor;
/**
* @className: AsyncEventBus
* @description: 异步 EventBus
* @date: 2022/5/13
* @author: cakin
*/
public class AsyncEventBus extends EventBus {
AsyncEventBus(String busName, EventExceptionHandler exceptionHandler, ThreadPoolExecutor executor) {
super(busName, exceptionHandler, executor);
}
AsyncEventBus(String busName, ThreadPoolExecutor executor) {
this(busName, null, executor);
}
AsyncEventBus(ThreadPoolExecutor executor) {
this("default_async", null, executor);
}
AsyncEventBus(EventExceptionHandler exceptionHandler, ThreadPoolExecutor executor) {
this("default_async", exceptionHandler, executor);
}
}
package concurrent.eventbus;
import java.lang.reflect.Method;
import java.lang.reflect.Modifier;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentLinkedQueue;
/**
* @className: Registry
* @description: 注册表维护了 topic 和 subscriber 之间的关系,当有 Event 被 post 之后,Dispatcher 需要知道该消息应该发送哪个 Subscriber 的实例和对应的方法
* @date: 2022/5/13
* @author: cakin
*/
public class Registry {
// 存储 Subscriber 集合和 topic 之间关系的 map
private final ConcurrentHashMap<String, ConcurrentLinkedQueue<Subscriber>> subscriberContainer = new ConcurrentHashMap<>();
public void bind(Object subscriber) {
// 获取 Subscriber Object 的方法集合,然后进行绑定
List<Method> subscribeMethods = getSubscribeMethods(subscriber);
subscribeMethods.forEach(m -> tierSubscriber(subscriber, m));
}
public void unbind(Object subscriber) {
// unbind 为了提高速度,只对 Subscriber 进行失效操作
subscriberContainer.forEach((key, queue) ->
queue.forEach(s -> {
if (s.getSubscribeObject() == subscriber) {
s.setDisable(true);
}
})
);
}
private void tierSubscriber(Object subscriber, Method method) {
final Subscribe subscribe = method.getDeclaredAnnotation(Subscribe.class);
String topic = subscribe.topic();
// 当某个 topic 没有 Subscriber Queue 的时候创建一个
subscriberContainer.computeIfAbsent(topic, key -> new ConcurrentLinkedQueue<>());
// 创建一个 subscriber 并且加入 subscriber 列表中
subscriberContainer.get(topic).add(new Subscriber(subscriber, method));
}
public ConcurrentLinkedQueue<Subscriber> scanSubscriber(final String topic) {
return subscriberContainer.get(topic);
}
private List<Method> getSubscribeMethods(Object subcriber) {
final List<Method> methods = new ArrayList<>();
Class<?> temp = subcriber.getClass();
// 不断获取所有的方法
while (temp != null) {
// 获取所有的方法
Method[] declaredMethods = temp.getDeclaredMethods();
// 只有 public 方法 && 有一个入参 && 被 @Subscribe 标记的方法才符合回调方法
Arrays.stream(declaredMethods)
.filter(m -> m.isAnnotationPresent(Subscribe.class)
&& m.getParameterCount() == 1
&& m.getModifiers() == Modifier.PUBLIC)
.forEach(methods::add);
temp = temp.getSuperclass();
}
return methods;
}
}
package concurrent.eventbus;
import java.lang.reflect.Method;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.Executor;
import java.util.concurrent.ExecutorService;
/**
* @className: Dispatcher
* @description: Event 广播 Dispatch
* @date: 2022/5/13
* @author: cakin
*/
public class Dispatcher {
private final Executor executorService;
private final EventExceptionHandler exceptionHandler;
public static final Executor SEQ_EXECUTOR_SERVICE = SeqExecutorService.INSTANCE;
public static final Executor PRE_THREAD_EXECUTOR_SERVICE = PreThreadExecutorService.INSTANCE;
public Dispatcher(Executor executorService, EventExceptionHandler exceptionHandler) {
this.executorService = executorService;
this.exceptionHandler = exceptionHandler;
}
public void dispatch(Bus bus, Registry registry, Object event, String topic) {
// 根据 topic 获取所有的 Subscriber 列表
ConcurrentLinkedQueue<Subscriber> subscribers = registry.scanSubscriber(topic);
if (null == subscribers) {
if (exceptionHandler != null) {
exceptionHandler.handler(new IllegalArgumentException("The topic" + topic + " note bind yet"), new BaseEventContext(bus.getBusName(), null, event));
return;
}
}
// 遍历所有的方法,并且通过反射的方式进行方法调用
subscribers.stream()
.filter(subscriber -> !subscriber.isDisable())
.filter(subscriber -> {
Method subcribeMethod = subscriber.getSubscribeMethod();
Class<?> aClass = subcribeMethod.getParameterTypes()[0];
return aClass.isAssignableFrom(event.getClass());
}).forEach(subscriber -> realInvokeSubscribe(subscriber, event, bus));
}
private void realInvokeSubscribe(Subscriber subscriber, Object event, Bus bus) {
Method subscribeMethod = subscriber.getSubscribeMethod();
Object subscribeObject = subscriber.getSubscribeObject();
executorService.execute(() -> {
try {
subscribeMethod.invoke(subscribeObject, event);
} catch (Exception e) {
if (null != exceptionHandler) {
exceptionHandler.handler(e, new BaseEventContext(bus.getBusName(), subscriber, event));
}
}
});
}
public void close() {
if (executorService instanceof ExecutorService) {
((ExecutorService) executorService).shutdown();
}
}
static Dispatcher newDispatcher(EventExceptionHandler exceptionHandler, Executor executor) {
return new Dispatcher(executor, exceptionHandler);
}
static Dispatcher seqDispatcher(EventExceptionHandler exceptionHandler) {
return new Dispatcher(SEQ_EXECUTOR_SERVICE, exceptionHandler);
}
static Dispatcher perThreadDispatcher(EventExceptionHandler exceptionHandler) {
return new Dispatcher(PRE_THREAD_EXECUTOR_SERVICE, exceptionHandler);
}
// 顺序执行的 ExecutorService
private static class SeqExecutorService implements Executor {
private final static SeqExecutorService INSTANCE = new SeqExecutorService();
@Override
public void execute(Runnable command) {
command.run();
}
}
// 每个线程负责一次消息推送
private static class PreThreadExecutorService implements Executor {
private final static PreThreadExecutorService INSTANCE = new PreThreadExecutorService();
@Override
public void execute(Runnable command) {
new Thread(command).start();
}
}
// 默认 EventContext 实现
private static class BaseEventContext implements EventContext {
private final String eventBusName;
private final Subscriber subscriber;
private final Object event;
private BaseEventContext(String eventBusName, Subscriber subscriber, Object event) {
this.eventBusName = eventBusName;
this.subscriber = subscriber;
this.event = event;
}
@Override
public String getSource() {
return this.eventBusName;
}
@Override
public Object getSubscriber() {
return subscriber != null ? subscriber.getSubscribeObject() : null;
}
@Override
public Method getSubscribe() {
return subscriber != null ? subscriber.getSubscribeMethod() : null;
}
@Override
public Object getEvent() {
return this.event;
}
}
}
package concurrent.eventbus;
import java.lang.reflect.Method;
/**
* @className: Subscriber
* @description: 封装了对象实例和被 @Subscribe 标记的方法
* @date: 2022/5/13
* @author: cakin
*/
public class Subscriber {
private final Object subscribeObject;
private final Method subscribeMethod;
private boolean disable = false;
public Subscriber(Object subscribeObject, Method subscribeMethod) {
this.subscribeObject = subscribeObject;
this.subscribeMethod = subscribeMethod;
}
public Object getSubscribeObject() {
return subscribeObject;
}
public Method getSubscribeMethod() {
return subscribeMethod;
}
public boolean isDisable() {
return disable;
}
public void setDisable(boolean disable) {
this.disable = disable;
}
}
package concurrent.eventbus;
/**
* @className: EventExceptionHandler
* @description: 异常处理
* @date: 2022/5/13
* @author: cakin
*/
public interface EventExceptionHandler {
void handler(Throwable cause, EventContext context);
}
package concurrent.eventbus;
import java.lang.reflect.Method;
/**
* @className: EventContext
* @description: 事件上下文
* @date: 2022/5/13
* @author: cakin
*/
public interface EventContext {
String getSource();
Object getSubscriber();
Method getSubscribe();
Object getEvent();
}
package concurrent.eventbus;
import java.lang.annotation.ElementType;
import java.lang.annotation.Retention;
import java.lang.annotation.RetentionPolicy;
import java.lang.annotation.Target;
/**
* @className: Subscribe
* @description: 注解在类的方法上,注解时可指定 topic,不指定的情况下未默认的 topic(default-topic)
* @date: 2022/5/13
* @author: cakin
*/
@Retention(RetentionPolicy.RUNTIME)
@Target(ElementType.METHOD)
public @interface Subscribe {
String topic() default "default-topic";
}
package concurrent.eventbus;
/**
* @className: SimpleSubsciber1
* @description: TODO
* @date: 2022/5/13
* @author: cakin
*/
public class SimpleSubscriber1 {
@Subscribe
public void method1(String message){
System.out.println("==SimpleSubscriber1==method1=="+message);
}
@Subscribe(topic = "test")
public void method2(String message){
System.out.println("==SimpleSubscriber1==method2=="+message);
}
}
package concurrent.eventbus;
/**
* @className: SimpleSubsciber2
* @description: SimpleSubsciber2
* @date: 2022/5/13
* @author: 贝医
*/
public class SimpleSubscriber2 {
@Subscribe
public void method1(String message){
System.out.println("==SimpleSubscriber2==method1=="+message);
}
@Subscribe(topic = "test")
public void method2(String message){
System.out.println("==SimpleSubscriber2==method2=="+message);
}
}
package concurrent.eventbus;
public class SyncTest {
public static void main(String[] args) {
Bus bus = new EventBus("Test");
bus.register(new SimpleSubscriber1());
bus.register(new SimpleSubscriber2());
bus.post("Hello");
System.out.println("---------");
bus.post("Hello", "test");
}
}
package concurrent.eventbus;
import java.util.concurrent.Executors;
import java.util.concurrent.ThreadPoolExecutor;
public class ASyncTest {
public static void main(String[] args) {
Bus bus = new AsyncEventBus("Test", (ThreadPoolExecutor) Executors.newFixedThreadPool(10));
bus.register(new SimpleSubscriber1());
bus.register(new SimpleSubscriber2());
bus.post("Hello");
System.out.println("---------");
bus.post("Hello", "test");
}
}
==SimpleSubscriber1==method1==Hello
==SimpleSubscriber2==method1==Hello
==SimpleSubscriber1==method2==Hello
==SimpleSubscriber2==method2==Hello
==SimpleSubscriber2==method1==Hello
==SimpleSubscriber1==method1==Hello
==SimpleSubscriber1==method2==Hello
==SimpleSubscriber2==method2==Hello
当检测鼠标x和y坐标时,最好像这样使用event.clientX和event.clientY: function show_coords(event){ var x=event.clientX;
我有以下代码: document.oncontextmenu = function(evt) { evt = evt || window.event; console.log(evt.
对于另一个问题,我遇到了一个似乎偶尔出现在 SO 的误解。一些提问者似乎认为触发器之于数据库就像事件之于 OOP 一样。 有没有人有一个很好的类比来解释为什么这是一个有缺陷的比较,以及误用它的后果?
$('body').keypress(function(event){ if(event.keyCode == 46){console.log('Delete Key Pressed')};
我正在制作一个“流体”文本区域,它根据内容调整其高度。我实际上正在尝试实现 this脚本。我有以下代码:https://ellie-app.com/Vjtvm6yrKWa1/4 问题是,当增加高度时,
我使用 Raphael .mouseover() 和 .mouseout() 事件来突出显示 SVG 中的某些元素。这工作正常,但在我单击一个元素后,我希望它停止突出显示。 在 Raphael doc
我目前正在开发一个应用程序,允许人们为在线广播电台安排“节目”。 我希望用户能够设置重复事件,例如:- “躁狂星期一”节目 - 每周一 9 点至 11 点“月中疯狂” - 每个月的第二个星期四“本月新
我有以下三个表格(简化版本): 已加载关卡: id(整数、主键、自动增量) globalId(整数,键) 日期(日期时间、键) serverId(Int,键) gamemodeId(Int,Key)
在我阅读 Gevent Tutorial 之后,我有一个关于 gevent.event.Event 的问题。 Event.set() 是否会唤醒所有被 Event.wait() 阻塞的函数? 就像下面
我对 cakephp ver3.1.3 没有经验 我按照说明实现了登录认证功能; http://book.cakephp.org/3.0/en/tutorials-and-examples/blog-
现在,我发送 10 个事件,每个事件有 1 个属性。但是当我想过滤特定事件并按属性选择事件时,在“事件属性”过滤器中仅显示前 7 个事件,而我为其余事件添加的事件仅显示“第一次”过滤器,为什么? 最佳
我不知道我的 Firefox 发生了什么! 我的aspx和javascript代码是这样的: function a() { alert(
中有3个事件fns重装 ,我可以对两者做同样的事情 reg-event-db和 reg-event-fx . reg-event-db之间的主要区别是什么, reg-event-fx和 reg-eve
我遇到了 Firefox keydown 行为,因为在没有聚焦于特定字段的情况下按下 Enter 键(实际上是任何键)不会触发 keydown 事件只会触发`按键事件。 这可能会非常令人困惑,因为 k
这是我的代码片段 public class Notation : INotifyPropertyChanged { public event PropertyChangedEventHandl
我可以在一个 Jsf2 xhtml 文件中有多个标签吗? 在那种情况下,关联的监听器将以什么顺序被调用? Mojarra 2.1.1/Apache Tomcat 7.0.22/PrimeFaces 3
我可以在一个 Jsf2 xhtml 文件中有多个标签吗? 在那种情况下,关联的监听器将以什么顺序被调用? Mojarra 2.1.1/Apache Tomcat 7.0.22/PrimeFaces 3
我有以下 JavaScript: $('#ge-display').click(function (event) { window.open('/googleearth/ge-display.ph
我需要确定触发事件的元素。 使用 event.target 获取相应的元素。 我可以从那里使用哪些属性? 引用 编号 节点名 我找不到关于它的大量信息,即使在 jQuery 上也是如此页,所以希望有人
我在pyGame中创建了一个Asteroidz克隆,并在pygame.vent.get()循环中有两个for Event,一个用于检查退出请求,以及游戏是否应该通过按空格键开始,然后在游戏中进一步尝试
我是一名优秀的程序员,十分优秀!