From 2887d365901ae5ea6cdf3e2f573cdfe0bc611442 Mon Sep 17 00:00:00 2001 From: tangmingyou <234767776@qq.com> Date: Wed, 13 Apr 2022 18:07:11 +0800 Subject: [PATCH] delay task queue schedule --- .../registry/ProtoMessageHandlerRegistry.java | 3 +- .../net/sopod/soim/core/session/NetUser.java | 16 ++++ .../soim/entry/config/EntryServerConfig.java | 2 +- .../entry/delay/NetUserDelayTaskManager.java | 95 +++++++++++++++++++ .../soim/entry/handler/HelloHandler.java | 3 + .../soim/entry/registry/DelayQueueTest.java | 64 +++++++++++++ .../entry/server/InboundImMessageHandler.java | 8 +- .../entry/server/ProtoMessageDispatcher.java | 25 +++++ .../net/sopod/soim/entry/worker/Worker.java | 2 +- .../sopod/soim/entry/worker/WorkerGroup.java | 10 -- .../src/main/resources/application.yml | 6 +- pom.xml | 16 ++++ 12 files changed, 229 insertions(+), 21 deletions(-) create mode 100644 im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java create mode 100644 im-entry/src/main/java/net/sopod/soim/entry/registry/DelayQueueTest.java create mode 100644 im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java 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 20defa8..901dc7d 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 @@ -1,5 +1,6 @@ package net.sopod.soim.core.registry; +import com.google.protobuf.MessageLite; import net.sopod.soim.common.util.ImClock; import net.sopod.soim.common.util.Reflects; import net.sopod.soim.core.handler.MessageHandler; @@ -62,7 +63,7 @@ public class ProtoMessageHandlerRegistry { } @Nullable - public static MessageHandler getTypeHandler(Class type) { + public static MessageHandler getTypeHandler(Class type) { if (CONTEXT_AWARE_AWAIT.getCount() > 0) { try { CONTEXT_AWARE_AWAIT.await(); diff --git a/im-core/src/main/java/net/sopod/soim/core/session/NetUser.java b/im-core/src/main/java/net/sopod/soim/core/session/NetUser.java index 89f7295..6337a4a 100644 --- a/im-core/src/main/java/net/sopod/soim/core/session/NetUser.java +++ b/im-core/src/main/java/net/sopod/soim/core/session/NetUser.java @@ -4,6 +4,7 @@ import io.netty.channel.Channel; import io.netty.util.AttributeKey; import java.lang.ref.WeakReference; +import java.util.concurrent.atomic.AtomicBoolean; public class NetUser { @@ -12,6 +13,8 @@ public class NetUser { private final WeakReference channel; + private final AtomicBoolean isActive = new AtomicBoolean(true); + public NetUser(Channel channel) { this.channel = new WeakReference<>(channel); } @@ -20,6 +23,19 @@ public class NetUser { return false; } + public boolean isActiveChannel() { + Channel chan = this.channel.get(); + return chan != null && chan.isActive(); + } + + public boolean isActive() { + return this.isActive.get(); + } + + public void inactive() { + this.isActive.set(false); + } + public void write(Object message) { write(message, false); } 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 index c289970..1336363 100644 --- 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 @@ -25,7 +25,7 @@ public class EntryServerConfig { private Integer port = 8088; /** 消息消费者线程数 */ - private Integer workerSize = 2; + private Integer workerSize; public Integer getWorkerSize() { if (workerSize == null || workerSize < 1) { diff --git a/im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java b/im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java new file mode 100644 index 0000000..8eb7d62 --- /dev/null +++ b/im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java @@ -0,0 +1,95 @@ +package net.sopod.soim.entry.delay; + +import com.google.common.base.Preconditions; +import com.google.protobuf.MessageLite; +import net.sopod.soim.common.util.ImClock; +import net.sopod.soim.core.handler.MessageHandler; +import net.sopod.soim.core.registry.ProtoMessageHandlerRegistry; +import net.sopod.soim.core.session.NetUser; +import net.sopod.soim.entry.server.ProtoMessageDispatcher; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.lang.ref.WeakReference; +import java.util.concurrent.*; + +/** + * NetUserDelayTaskManager + * 延时任务执行 + * + * @author tmy + * @date 2022-04-13 16:59 + */ +public class NetUserDelayTaskManager { + + private static final Logger logger = LoggerFactory.getLogger(NetUserDelayTaskManager.class); + + private static final DelayQueue DELAY_QUEUE; + + static { + DELAY_QUEUE = new DelayQueue<>(); + ScheduledExecutorService scheduledExecutor = Executors.newSingleThreadScheduledExecutor(); + scheduledExecutor.scheduleWithFixedDelay(() -> { + try { + DelayTask delayTask; + // 连续批量执行50个任务 + for (int i = 0; i < 50; i++) { + delayTask = DELAY_QUEUE.poll(); + if (delayTask == null) { + break; + } + NetUser netUser = delayTask.getNetUser(); + if (netUser != null && netUser.isActive()) { + ProtoMessageDispatcher.dispatch(netUser, delayTask.getTaskMsg()); + } else { + logger.info("delayTask netUser inactive, task cancel!"); + } + } + }catch (Throwable e) { + logger.error("netUser delay task schedule error: ", e); + } + }, 100,100, TimeUnit.MILLISECONDS); + } + + /** + * TODO boolean netUser inactive cancel + */ + public static void addTask(NetUser netUser, MessageLite taskMsg, long delay, TimeUnit unit) { + Preconditions.checkNotNull(taskMsg); + Preconditions.checkNotNull(unit); + Preconditions.checkArgument(delay > 0, "延迟时间值必须大于0"); + // 检查任务消息有无对应 handler 执行 + MessageHandler typeHandler = ProtoMessageHandlerRegistry.getTypeHandler(taskMsg.getClass()); + if (typeHandler == null) { + throw new IllegalStateException("not found taskMsg handler for type" + taskMsg.getClass() + "!"); + } + DELAY_QUEUE.add(new DelayTask(netUser, taskMsg, ImClock.millis() + unit.toMillis(delay))); + } + + private static class DelayTask implements Delayed { + // 如果 netUser 执行时不存在了,任务取消 + private final WeakReference netUser; + private final MessageLite taskMsg; + private final long time; + public DelayTask(NetUser netUser, MessageLite taskMsg, long time) { + this.netUser = new WeakReference<>(netUser); + this.taskMsg = taskMsg; + this.time = time; + } + @Override + public long getDelay(TimeUnit timeUnit) { + return time - ImClock.millis(); + } + @Override + public int compareTo(Delayed delayed) { + return Long.compare(time, delayed.getDelay(TimeUnit.MILLISECONDS)); + } + public MessageLite getTaskMsg() { + return this.taskMsg; + } + public NetUser getNetUser() { + return netUser.get(); + } + } + +} 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 4ef6330..0574d91 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,6 +4,7 @@ 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.entry.delay.NetUserDelayTaskManager; import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator; import net.sopod.soim.logic.user.service.UserService; import org.apache.dubbo.config.annotation.DubboReference; @@ -15,6 +16,7 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.stereotype.Service; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; /** * HelloHandler @@ -46,6 +48,7 @@ public class HelloHandler extends NetUserMessageHandler { RpcContext.getServiceContext().getCompletableFuture().whenComplete((res, err) -> { logger.info("async sayHello, {}", res); }); + NetUserDelayTaskManager.addTask(netUser, msg, 5, TimeUnit.SECONDS); return null; } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/registry/DelayQueueTest.java b/im-entry/src/main/java/net/sopod/soim/entry/registry/DelayQueueTest.java new file mode 100644 index 0000000..9282fc6 --- /dev/null +++ b/im-entry/src/main/java/net/sopod/soim/entry/registry/DelayQueueTest.java @@ -0,0 +1,64 @@ +package net.sopod.soim.entry.registry; + +import net.sopod.soim.common.util.ImClock; + +import java.util.concurrent.DelayQueue; +import java.util.concurrent.Delayed; +import java.util.concurrent.TimeUnit; + +/** + * DelayQueueTest + * + * @author tmy + * @date 2022-04-13 16:41 + */ +public class DelayQueueTest { + + public static class DelayEle implements Delayed { + + private final TimeUnit unit; + + private final long time; + + public DelayEle(long delay, TimeUnit unit) { + this.unit = unit; + this.time = ImClock.millis() + unit.toMillis(delay); + } + + @Override + public long getDelay(TimeUnit timeUnit) { + System.out.println("[get delay]..." + this.time); + return time - ImClock.millis(); + } + + @Override + public int compareTo(Delayed delayed) { + return Long.compare(this.time, delayed.getDelay(unit)); + } + + public long getTime() { + return this.time; + } + + } + + public static void main(String[] args) throws InterruptedException { + DelayQueue delayQueue = new DelayQueue<>(); + DelayEle ele1 = new DelayEle(10, TimeUnit.SECONDS); + DelayEle ele2 = new DelayEle(15, TimeUnit.SECONDS); + DelayEle ele3 = new DelayEle(20, TimeUnit.SECONDS); + delayQueue.add(ele1); + delayQueue.add(ele2); + delayQueue.add(ele3); + for (int i = 0; i < 3;) { + DelayEle ele = (DelayEle)delayQueue.poll(); + if (ele != null) { + System.out.println(ele.getTime()); + i ++; + } + Thread.sleep(1000); + } + + } + +} diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java index cb16a5a..fc7b41c 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java @@ -38,6 +38,8 @@ public class InboundImMessageHandler extends SimpleChannelInboundHandler netUser = ctx.channel().attr(NetUser.NET_USER_KEY); + netUser.get().inactive(); ctx.fireChannelInactive(); } @@ -47,11 +49,7 @@ public class InboundImMessageHandler extends SimpleChannelInboundHandler netUserAttr = ctx.channel().attr(NetUser.NET_USER_KEY); NetUser netUser = netUserAttr.get(); // TODO dispatcher - MessageHandler typeHandler = (MessageHandler) ProtoMessageHandlerRegistry - .getTypeHandler(messageLite.getClass()); - - Worker worker = WorkerGroup.next(); - worker.execute(() -> { typeHandler.exec(netUser, messageLite); }); + ProtoMessageDispatcher.dispatch(netUser, messageLite); } } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java b/im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java new file mode 100644 index 0000000..8a8c3c0 --- /dev/null +++ b/im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java @@ -0,0 +1,25 @@ +package net.sopod.soim.entry.server; + +import com.google.protobuf.MessageLite; +import net.sopod.soim.core.handler.MessageHandler; +import net.sopod.soim.core.registry.ProtoMessageHandlerRegistry; +import net.sopod.soim.core.session.NetUser; +import net.sopod.soim.entry.worker.Worker; +import net.sopod.soim.entry.worker.WorkerGroup; + +/** + * ProtoMessageDispatcher + * + * @author tmy + * @date 2022-04-13 17:21 + */ +public class ProtoMessageDispatcher { + + public static void dispatch(NetUser netUser, MessageLite message) { + MessageHandler typeHandler = ProtoMessageHandlerRegistry + .getTypeHandler(message.getClass()); + Worker worker = WorkerGroup.next(); + worker.execute(() -> { typeHandler.exec(netUser, message); }); + } + +} diff --git a/im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java b/im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java index 1d8b132..b0aceab 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java @@ -31,7 +31,7 @@ public class Worker implements EventHandler, EventFactory ); this.disruptor = new Disruptor<>( this, - 8 * 1024, + 16 * 1024, executor, ProducerType.SINGLE, new BlockingWaitStrategy()); diff --git a/im-entry/src/main/java/net/sopod/soim/entry/worker/WorkerGroup.java b/im-entry/src/main/java/net/sopod/soim/entry/worker/WorkerGroup.java index 5526e79..66c857f 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/worker/WorkerGroup.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/worker/WorkerGroup.java @@ -39,16 +39,6 @@ public class WorkerGroup { return WORKERS[counter.incrementAndGet() % WORKERS.length]; } - public static void publish(NetUser netUser, GeneratedMessageV3 message) { - next().execute(() -> { - - }); - } - - public static void publish(Account netUser, GeneratedMessageV3 message) { - - } - public static void shutdown() { for (Worker worker : WORKERS) { worker.shutdown(); 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 ca793f1..4c75ab8 100644 --- a/im-logic/im-segment-id/src/main/resources/application.yml +++ b/im-logic/im-segment-id/src/main/resources/application.yml @@ -10,11 +10,11 @@ spring: minimum-idle: 1 maximum-pool-size: 8 connection-timeout: 2000 - idle-timeout: 600000 # 10分钟空闲关闭 - max-lifetime: 1200000 # 20分钟最大存活时间 + idle-timeout: 300000 # 5分钟空闲关闭 + max-lifetime: 600000 # 10分钟最大存活时间 validation-timeout: 2000 connection-init-sql: select 1 - keepalive-time: 60000 # 连接存活时间,小于maxLifetime, 最小30秒, 空闲30秒后移除连接测试通过再添加回池 + keepalive-time: 30000 # 连接存活时间,小于maxLifetime, 最小30秒, 空闲30秒后移除连接测试通过再添加回池 redis: host: 124.222.131.236 port: 3379 diff --git a/pom.xml b/pom.xml index 7aaa726..2e4cc26 100644 --- a/pom.xml +++ b/pom.xml @@ -277,4 +277,20 @@ + + + aliyun + aliyun + https://maven.aliyun.com/repository/public + + + + + + aliyun + aliyun + https://maven.aliyun.com/repository/public + + + \ No newline at end of file