Browse Source

delay task queue schedule

master
tangmingyou 4 years ago
parent
commit
2887d36590
  1. 3
      im-core/src/main/java/net/sopod/soim/core/registry/ProtoMessageHandlerRegistry.java
  2. 16
      im-core/src/main/java/net/sopod/soim/core/session/NetUser.java
  3. 2
      im-entry/src/main/java/net/sopod/soim/entry/config/EntryServerConfig.java
  4. 95
      im-entry/src/main/java/net/sopod/soim/entry/delay/NetUserDelayTaskManager.java
  5. 3
      im-entry/src/main/java/net/sopod/soim/entry/handler/HelloHandler.java
  6. 64
      im-entry/src/main/java/net/sopod/soim/entry/registry/DelayQueueTest.java
  7. 8
      im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java
  8. 25
      im-entry/src/main/java/net/sopod/soim/entry/server/ProtoMessageDispatcher.java
  9. 2
      im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java
  10. 10
      im-entry/src/main/java/net/sopod/soim/entry/worker/WorkerGroup.java
  11. 6
      im-logic/im-segment-id/src/main/resources/application.yml
  12. 16
      pom.xml

3
im-core/src/main/java/net/sopod/soim/core/registry/ProtoMessageHandlerRegistry.java

@ -1,5 +1,6 @@
package net.sopod.soim.core.registry; package net.sopod.soim.core.registry;
import com.google.protobuf.MessageLite;
import net.sopod.soim.common.util.ImClock; import net.sopod.soim.common.util.ImClock;
import net.sopod.soim.common.util.Reflects; import net.sopod.soim.common.util.Reflects;
import net.sopod.soim.core.handler.MessageHandler; import net.sopod.soim.core.handler.MessageHandler;
@ -62,7 +63,7 @@ public class ProtoMessageHandlerRegistry {
} }
@Nullable @Nullable
public static <T> MessageHandler<T> getTypeHandler(Class<T> type) { public static <T> MessageHandler<T> getTypeHandler(Class<? extends MessageLite> type) {
if (CONTEXT_AWARE_AWAIT.getCount() > 0) { if (CONTEXT_AWARE_AWAIT.getCount() > 0) {
try { try {
CONTEXT_AWARE_AWAIT.await(); CONTEXT_AWARE_AWAIT.await();

16
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 io.netty.util.AttributeKey;
import java.lang.ref.WeakReference; import java.lang.ref.WeakReference;
import java.util.concurrent.atomic.AtomicBoolean;
public class NetUser { public class NetUser {
@ -12,6 +13,8 @@ public class NetUser {
private final WeakReference<Channel> channel; private final WeakReference<Channel> channel;
private final AtomicBoolean isActive = new AtomicBoolean(true);
public NetUser(Channel channel) { public NetUser(Channel channel) {
this.channel = new WeakReference<>(channel); this.channel = new WeakReference<>(channel);
} }
@ -20,6 +23,19 @@ public class NetUser {
return false; 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) { public void write(Object message) {
write(message, false); write(message, false);
} }

2
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 port = 8088;
/** 消息消费者线程数 */ /** 消息消费者线程数 */
private Integer workerSize = 2; private Integer workerSize;
public Integer getWorkerSize() { public Integer getWorkerSize() {
if (workerSize == null || workerSize < 1) { if (workerSize == null || workerSize < 1) {

95
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<DelayTask> 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<MessageLite> 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> 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();
}
}
}

3
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.handler.NetUserMessageHandler;
import net.sopod.soim.core.session.NetUser; import net.sopod.soim.core.session.NetUser;
import net.sopod.soim.data.msg.hello.HelloPB; 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.api.segmentid.core.SegmentIdGenerator;
import net.sopod.soim.logic.user.service.UserService; import net.sopod.soim.logic.user.service.UserService;
import org.apache.dubbo.config.annotation.DubboReference; import org.apache.dubbo.config.annotation.DubboReference;
@ -15,6 +16,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
/** /**
* HelloHandler * HelloHandler
@ -46,6 +48,7 @@ public class HelloHandler extends NetUserMessageHandler<HelloPB.Hello> {
RpcContext.getServiceContext().getCompletableFuture().whenComplete((res, err) -> { RpcContext.getServiceContext().getCompletableFuture().whenComplete((res, err) -> {
logger.info("async sayHello, {}", res); logger.info("async sayHello, {}", res);
}); });
NetUserDelayTaskManager.addTask(netUser, msg, 5, TimeUnit.SECONDS);
return null; return null;
} }

64
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<Delayed> 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);
}
}
}

8
im-entry/src/main/java/net/sopod/soim/entry/server/InboundImMessageHandler.java

@ -38,6 +38,8 @@ public class InboundImMessageHandler extends SimpleChannelInboundHandler<Message
@Override @Override
public void channelInactive(ChannelHandlerContext ctx) { public void channelInactive(ChannelHandlerContext ctx) {
Attribute<NetUser> netUser = ctx.channel().attr(NetUser.NET_USER_KEY);
netUser.get().inactive();
ctx.fireChannelInactive(); ctx.fireChannelInactive();
} }
@ -47,11 +49,7 @@ public class InboundImMessageHandler extends SimpleChannelInboundHandler<Message
Attribute<NetUser> netUserAttr = ctx.channel().attr(NetUser.NET_USER_KEY); Attribute<NetUser> netUserAttr = ctx.channel().attr(NetUser.NET_USER_KEY);
NetUser netUser = netUserAttr.get(); NetUser netUser = netUserAttr.get();
// TODO dispatcher // TODO dispatcher
MessageHandler<MessageLite> typeHandler = (MessageHandler<MessageLite>) ProtoMessageHandlerRegistry ProtoMessageDispatcher.dispatch(netUser, messageLite);
.getTypeHandler(messageLite.getClass());
Worker worker = WorkerGroup.next();
worker.execute(() -> { typeHandler.exec(netUser, messageLite); });
} }
} }

25
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<MessageLite> typeHandler = ProtoMessageHandlerRegistry
.getTypeHandler(message.getClass());
Worker worker = WorkerGroup.next();
worker.execute(() -> { typeHandler.exec(netUser, message); });
}
}

2
im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java

@ -31,7 +31,7 @@ public class Worker implements EventHandler<TaskEvent>, EventFactory<TaskEvent>
); );
this.disruptor = new Disruptor<>( this.disruptor = new Disruptor<>(
this, this,
8 * 1024, 16 * 1024,
executor, executor,
ProducerType.SINGLE, ProducerType.SINGLE,
new BlockingWaitStrategy()); new BlockingWaitStrategy());

10
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]; 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() { public static void shutdown() {
for (Worker worker : WORKERS) { for (Worker worker : WORKERS) {
worker.shutdown(); worker.shutdown();

6
im-logic/im-segment-id/src/main/resources/application.yml

@ -10,11 +10,11 @@ spring:
minimum-idle: 1 minimum-idle: 1
maximum-pool-size: 8 maximum-pool-size: 8
connection-timeout: 2000 connection-timeout: 2000
idle-timeout: 600000 # 10分钟空闲关闭 idle-timeout: 300000 # 5分钟空闲关闭
max-lifetime: 1200000 # 20分钟最大存活时间 max-lifetime: 600000 # 10分钟最大存活时间
validation-timeout: 2000 validation-timeout: 2000
connection-init-sql: select 1 connection-init-sql: select 1
keepalive-time: 60000 # 连接存活时间,小于maxLifetime, 最小30秒, 空闲30秒后移除连接测试通过再添加回池 keepalive-time: 30000 # 连接存活时间,小于maxLifetime, 最小30秒, 空闲30秒后移除连接测试通过再添加回池
redis: redis:
host: 124.222.131.236 host: 124.222.131.236
port: 3379 port: 3379

16
pom.xml

@ -277,4 +277,20 @@
</plugins> </plugins>
</build> </build>
<repositories>
<repository>
<id>aliyun</id>
<name>aliyun</name>
<url>https://maven.aliyun.com/repository/public</url>
</repository>
</repositories>
<pluginRepositories>
<pluginRepository>
<id>aliyun</id>
<name>aliyun</name>
<url>https://maven.aliyun.com/repository/public</url>
</pluginRepository>
</pluginRepositories>
</project> </project>
Loading…
Cancel
Save