AKKA Event Bus
事件机制就用于当前运行环境,与集群环境不同,详细见AKKA 集群中的发布与订阅Distributed Publish Subscribe in Cluster
简单实现示例
1package event 2 3import akka.actor.AbstractActor 4import akka.actor.ActorRef 5import akka.actor.ActorSystem 6import akka.actor.Props 7import akka.event.japi.LookupEventBus 8import akka.japi.pf.ReceiveBuilder 9import com.typesafe.config.ConfigFactory 10 11/** 12 * Created by: tankx 13 * Date: 2019/7/18 14 * Description: 事件与监听 15 */ 16object EventBus : LookupEventBus<MyEvent, ActorRef, String>() {//参数(事件类型,订阅者类型,用于区分事件定义的类型) 17 18 override fun classify(event: MyEvent): String {//用于区分不同事件(事件类型) 19 return event.type 20 } 21 22 override fun publish(event: MyEvent, subscriber: ActorRef) { 23 subscriber.tell(event, ActorRef.noSender()) 24 } 25 26 //期望的事件类型的数量 27 override fun mapSize(): Int { 28 return 1000 29 } 30 31 override fun compareSubscribers(a: ActorRef, b: ActorRef): Int { 32 return a.compareTo(b) 33 } 34 35 36} 37 38//订阅actor 39class SubActor : AbstractActor() { 40 41 override fun createReceive(): Receive { 42 return ReceiveBuilder.create().matchAny(this::receive).build() 43 } 44 45 fun receive(msg: Any) { 46 47 println("收到消息: $msg") 48 49 } 50 51} 52 53fun main() { 54 55 var system: ActorSystem = ActorSystem.create("system"); 56 57 var eventActor = system.actorOf(Props.create(SubActor::class.java)) 58 59 60 EventBus.subscribe(eventActor, "aaa")//(订阅者,事件类型) 61 EventBus.subscribe(eventActor, "bbb") 62 EventBus.subscribe(eventActor, "ccc") 63 EventBus.subscribe(eventActor, "ddd") 64 65 66 EventBus.publish(MyEvent("aaa", "数据")) 67 EventBus.publish(MyEvent("bbb", "数据")) 68 EventBus.publish(MyEvent("ccc", "数据")) 69 EventBus.publish(MyEvent("ccc", "数据")) 70}