diff --git a/README.md b/README.md
index 3c58a6e..a5dfffc 100644
--- a/README.md
+++ b/README.md
@@ -52,4 +52,5 @@ router <--> das
das shardingjdbc
table struct, sharding roles
logic biz
-
\ No newline at end of file
+
+请求响应消息队列异步处理,减少线程 cpu 占用, dubbo async
\ No newline at end of file
diff --git a/im-client/pom.xml b/im-client/pom.xml
index f8c7160..d3949d0 100644
--- a/im-client/pom.xml
+++ b/im-client/pom.xml
@@ -25,6 +25,11 @@
io.netty
netty-all
+
+ com.beust
+ jcommander
+ 1.82
+
\ No newline at end of file
diff --git a/im-client/src/main/java/net/sopod/soim/client/ClientMain.java b/im-client/src/main/java/net/sopod/soim/client/ClientMain.java
index f0c2dd1..27e5eef 100644
--- a/im-client/src/main/java/net/sopod/soim/client/ClientMain.java
+++ b/im-client/src/main/java/net/sopod/soim/client/ClientMain.java
@@ -4,6 +4,9 @@ import net.sopod.soim.client.net.ImNetClient;
/**
* Main
+ * jcommander:
+ * https://www.codenong.com/b-jcommander-parsing-command-line-parameters/
+ * https://www.mianshigee.com/project/jcommander
*
* @author tmy
* @date 2022-03-27 22:32
diff --git a/im-core/src/main/java/net/sopod/soim/core/registry/ProtoMessageHandlerRegistry.java b/im-core/src/main/java/net/sopod/soim/core/registry/ProtoMessageHandlerRegistry.java
index 86990ef..20defa8 100644
--- a/im-core/src/main/java/net/sopod/soim/core/registry/ProtoMessageHandlerRegistry.java
+++ b/im-core/src/main/java/net/sopod/soim/core/registry/ProtoMessageHandlerRegistry.java
@@ -47,6 +47,7 @@ public class ProtoMessageHandlerRegistry {
try {
Class> type = genericTypes.size() == 0 ? Object.class : Class.forName(genericTypes.get(0));
MessageHandler> existHandler = TYPE_HANDLER_MAP.putIfAbsent(type, handler);
+ logger.debug("registe msg type {} for handler {}", type, handler);
if (existHandler != null) {
// 消息类型有重复的 handler!
throw new IllegalStateException("msg type " + type + " handler duplicate; " +
diff --git a/im-entry/pom.xml b/im-entry/pom.xml
index fe0dd44..614562f 100644
--- a/im-entry/pom.xml
+++ b/im-entry/pom.xml
@@ -21,6 +21,16 @@
im-core
${soim.version}
+
+ net.sopod
+ im-segment-id-api
+ ${soim.version}
+
+
+ net.sopod
+ im-logic-user-api
+ ${soim.version}
+
org.projectlombok
lombok
@@ -62,6 +72,10 @@
org.apache.dubbo
dubbo-registry-nacos
+
+ com.lmax
+ disruptor
+
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/config/EntryServerConfig.java b/im-entry/src/main/java/net/sopod/soim/entry/config/EntryServerConfig.java
new file mode 100644
index 0000000..327f18d
--- /dev/null
+++ b/im-entry/src/main/java/net/sopod/soim/entry/config/EntryServerConfig.java
@@ -0,0 +1,31 @@
+package net.sopod.soim.entry.config;
+
+import lombok.Data;
+import org.springframework.boot.context.properties.ConfigurationProperties;
+import org.springframework.boot.context.properties.EnableConfigurationProperties;
+import org.springframework.stereotype.Component;
+
+/**
+ * EntryServerConfig
+ *
+ * @author tmy
+ * @date 2022-04-11 16:08
+ */
+@Component
+@ConfigurationProperties(prefix = "entry")
+@EnableConfigurationProperties
+@Data
+public class EntryServerConfig {
+
+ /** 消息消费者线程数 */
+ private Integer workerSize;
+
+ public Integer getWorkerSize() {
+ if (workerSize == null || workerSize < 1) {
+ workerSize = Runtime.getRuntime().availableProcessors();
+ workerSize = Math.max(workerSize, 4);
+ }
+ return workerSize;
+ }
+
+}
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/handler/HelloHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/handler/HelloHandler.java
index 09bf2df..4ef6330 100644
--- a/im-entry/src/main/java/net/sopod/soim/entry/handler/HelloHandler.java
+++ b/im-entry/src/main/java/net/sopod/soim/entry/handler/HelloHandler.java
@@ -4,8 +4,18 @@ import com.google.protobuf.MessageLite;
import net.sopod.soim.core.handler.NetUserMessageHandler;
import net.sopod.soim.core.session.NetUser;
import net.sopod.soim.data.msg.hello.HelloPB;
+import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator;
+import net.sopod.soim.logic.user.service.UserService;
+import org.apache.dubbo.config.annotation.DubboReference;
+import org.apache.dubbo.config.annotation.Method;
+import org.apache.dubbo.rpc.RpcContext;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
+import java.util.concurrent.CompletableFuture;
+
/**
* HelloHandler
*
@@ -15,11 +25,27 @@ import org.springframework.stereotype.Service;
@Service
public class HelloHandler extends NetUserMessageHandler {
+ private static final Logger logger = LoggerFactory.getLogger(HelloHandler.class);
+
+ @DubboReference(methods = {@Method(name = "sayHello", async = true)})
+ private UserService userService;
+
+ @Autowired
+ private SegmentIdGenerator segmentIdGenerator;
+
@Override
public MessageLite handle(NetUser netUser, HelloPB.Hello msg) {
- System.out.println("get hello message");
- System.out.println(msg);
- System.out.println(msg.getStr());
+ logger.info("get hello message: {}", segmentIdGenerator.nextId("im-entry-hello"));
+ logger.info("msg={}, {}", msg.getId(), msg.getStr());
+ CompletableFuture hi = userService.sayHi("黄绿");
+ hi.whenComplete((res, err) -> {
+ logger.info("async hi, {}", res);
+ });
+ String hello = userService.sayHello("lastJet");
+ logger.info("hello:{}", hello);
+ RpcContext.getServiceContext().getCompletableFuture().whenComplete((res, err) -> {
+ logger.info("async sayHello, {}", res);
+ });
return null;
}
diff --git a/im-entry/src/main/java/net/sopod/soim/entry/worker/DisruptorTest.java b/im-entry/src/main/java/net/sopod/soim/entry/worker/DisruptorTest.java
new file mode 100644
index 0000000..6e1483c
--- /dev/null
+++ b/im-entry/src/main/java/net/sopod/soim/entry/worker/DisruptorTest.java
@@ -0,0 +1,78 @@
+package net.sopod.soim.entry.worker;
+
+import com.lmax.disruptor.*;
+import com.lmax.disruptor.dsl.Disruptor;
+import com.lmax.disruptor.dsl.ProducerType;
+import io.netty.util.concurrent.DefaultThreadFactory;
+import net.sopod.soim.common.util.ImClock;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+/**
+ * Worker
+ *
+ * @author tmy
+ * @date 2022-04-11 17:05
+ */
+public class DisruptorTest {
+ public static class LongEvent {
+ int id;
+ private Long value;
+ public Long getValue() {
+ return value;
+ }
+ public void setValue(Long value) {
+ this.value = value;
+ }
+ }
+ public static class LongEventFactory implements EventFactory {
+ private static final AtomicInteger counter = new AtomicInteger(1);
+ @Override
+ public LongEvent newInstance() {
+ int count = counter.getAndIncrement();
+ System.out.println("instance:"+ count);
+ LongEvent event = new LongEvent();
+ event.id = count;
+ return event;
+ }
+ }
+ public static class LongEventHandler implements EventHandler {
+ @Override
+ public void onEvent(LongEvent event, long sequence, boolean endOfBatch) throws Exception {
+ System.out.println("消费者:"+event.id + "," + sequence + "," + event.getValue() + "," + Thread.currentThread().getName());
+ Thread.sleep(20);
+ // event值置空,回收内存
+ event.setValue(null);
+ }
+ }
+ public static void main(String[] args) throws InterruptedException, InsufficientCapacityException {
+ LongEventFactory eventFactory = new LongEventFactory();
+ int ringBufferSize = 128;//64 * 1024;
+ Disruptor disruptor = new Disruptor(
+ eventFactory,
+ ringBufferSize,
+ new DefaultThreadFactory("disruptor"),
+ ProducerType.SINGLE,
+ new BlockingWaitStrategy()
+ );
+ disruptor.handleEventsWith(new LongEventHandler());
+ disruptor.start();
+
+ // 7.创建RingBuffer容器
+ RingBuffer ringBuffer = disruptor.getRingBuffer();
+ for (long i = 0; i < 1000; i++) {
+ long start = ImClock.millis();
+ // 没有空闲位置时会阻塞
+ long sequence = ringBuffer.tryNext();
+ //long sequence = ringBuffer.next();
+ LongEvent longEvent = ringBuffer.get(sequence);
+ longEvent.setValue(i);
+ Thread.sleep(10);
+ ringBuffer.publish(sequence);
+ System.out.println("生产者:" + i + "," + (ImClock.millis() - start) + "ms");
+ }
+
+ Thread.sleep(5000);
+ disruptor.shutdown();
+ }
+}
diff --git a/im-entry/src/main/resources/application.yml b/im-entry/src/main/resources/application.yml
index ae7f1b0..44f4153 100644
--- a/im-entry/src/main/resources/application.yml
+++ b/im-entry/src/main/resources/application.yml
@@ -5,4 +5,4 @@ dubbo:
application:
name: ${spring.application.name}
registry:
- address: nacos://124.222.131.236:3848
+ address: nacos://124.222.131.236:3848
\ No newline at end of file
diff --git a/im-logic-api/im-logic-user-api/src/main/java/net/sopod/soim/logic/user/service/UserService.java b/im-logic-api/im-logic-user-api/src/main/java/net/sopod/soim/logic/user/service/UserService.java
index adc0ecf..4f3f0d5 100644
--- a/im-logic-api/im-logic-user-api/src/main/java/net/sopod/soim/logic/user/service/UserService.java
+++ b/im-logic-api/im-logic-user-api/src/main/java/net/sopod/soim/logic/user/service/UserService.java
@@ -1,5 +1,8 @@
package net.sopod.soim.logic.user.service;
+
+import java.util.concurrent.CompletableFuture;
+
/**
* UserService
*
@@ -9,8 +12,10 @@ package net.sopod.soim.logic.user.service;
public interface UserService {
/**
- *
+ * 异步接口测试
*/
- void a();
+ CompletableFuture sayHi(String name);
+
+ String sayHello(String name);
}
diff --git a/im-logic/im-logic-user/src/main/java/net/sopod/soim/logic/user/LogicUserApplication.java b/im-logic/im-logic-user/src/main/java/net/sopod/soim/logic/user/LogicUserApplication.java
index be34652..a77cd66 100644
--- a/im-logic/im-logic-user/src/main/java/net/sopod/soim/logic/user/LogicUserApplication.java
+++ b/im-logic/im-logic-user/src/main/java/net/sopod/soim/logic/user/LogicUserApplication.java
@@ -1,15 +1,21 @@
package net.sopod.soim.logic.user;
+import org.apache.dubbo.config.spring.context.annotation.EnableDubbo;
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.SpringBootApplication;
+
/**
* UserMain
*
* @author tmy
* @date 2022-03-26 01:41
*/
+@EnableDubbo(scanBasePackages = {"net.sopod.soim.logic.user.service"})
+@SpringBootApplication
public class LogicUserApplication {
public static void main(String[] args) {
-
+ SpringApplication.run(LogicUserApplication.class, args);
}
}
diff --git a/im-logic/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/UserServiceImpl.java b/im-logic/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/UserServiceImpl.java
index 4da2f84..c79842c 100644
--- a/im-logic/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/UserServiceImpl.java
+++ b/im-logic/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/UserServiceImpl.java
@@ -1,10 +1,26 @@
package net.sopod.soim.logic.user.service;
+import org.apache.dubbo.config.annotation.DubboService;
+
+import java.util.concurrent.CompletableFuture;
+
/**
* UserServiceImpl
*
* @author tmy
* @date 2022-04-04 10:03
*/
-public class UserServiceImpl {
+@DubboService
+public class UserServiceImpl implements UserService {
+
+ @Override
+ public CompletableFuture sayHi(String name) {
+ return CompletableFuture.completedFuture("hi," + name + "!");
+ }
+
+ @Override
+ public String sayHello(String name) {
+ return "hello," + name + "!";
+ }
+
}
diff --git a/im-logic/im-logic-user/src/main/resources/application.yml b/im-logic/im-logic-user/src/main/resources/application.yml
new file mode 100644
index 0000000..7f32fc9
--- /dev/null
+++ b/im-logic/im-logic-user/src/main/resources/application.yml
@@ -0,0 +1,11 @@
+spring:
+ application:
+ name: im-logic-user
+
+dubbo:
+ application:
+ name: ${spring.application.name}
+ registry:
+ address: nacos://124.222.131.236:3848
+ protocol:
+ port: 3002
diff --git a/im-logic/im-logic-user/src/main/resources/setting.yml b/im-logic/im-logic-user/src/main/resources/setting.yml
deleted file mode 100644
index e69de29..0000000
diff --git a/im-logic/im-segment-id/src/main/resources/application.yml b/im-logic/im-segment-id/src/main/resources/application.yml
index 48180f0..dcde657 100644
--- a/im-logic/im-segment-id/src/main/resources/application.yml
+++ b/im-logic/im-segment-id/src/main/resources/application.yml
@@ -6,6 +6,15 @@ spring:
url: jdbc:mysql://cd-cdb-mrz9fw80.sql.tencentcdb.com:61843/soim_db?serverTimezone=GMT%2B8
username: root
password: sopod@2347#
+ hikari:
+ minimum-idle: 1
+ maximum-pool-size: 8
+ connection-timeout: 2000
+ idle-timeout: 600000 # 10分钟空闲关闭
+ max-lifetime: 1200000 # 20分钟最大存活时间
+ validation-timeout: 2000
+ connection-init-sql: select 1
+ keepalive-time: 60000 # 连接存活时间,小于maxLifetime, 最小30秒, 空闲30秒后移除连接测试通过再添加回池
redis:
host: 124.222.131.236
port: 3379