From 17145b7ba286eb050679689e8178665d15ed7da6 Mon Sep 17 00:00:00 2001 From: tangmingyou <234767776@qq.com> Date: Mon, 23 May 2022 23:45:39 +0800 Subject: [PATCH] =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E5=A5=BD=E5=8F=8B,=E5=A5=BD?= =?UTF-8?q?=E5=8F=8B=E5=88=97=E8=A1=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../net/sopod/soim/client/cmd/CmdEnum.java | 5 ++ .../soim/client/cmd/args/ArgsAddFriend.java | 4 +- .../client/handler/cmd/AddFriendHandler.java | 16 +++++- .../client/handler/cmd/FriendsHandler.java | 27 ++++++++++ .../handler/msg/ResAddFriendHandler.java | 27 ++++++++++ .../handler/msg/ResFriendListHandler.java | 30 +++++++++++ .../handler/msg/ResFriendSearchHandler.java | 4 +- .../handler/msg/ResOnlineUserListHandler.java | 4 +- .../handler/msg/ResTokenAuthHandler.java | 4 +- .../client/handler/msg/TextChatHandler.java | 4 +- .../soim/client/session/MessageHandler.java | 2 +- .../soim/das/user/api/config/LogicTables.java | 2 + .../das/user/api/model/entity/ImFriend.java | 3 +- .../soim/das/user/service/FriendDasImpl.java | 11 ++++ .../src/main/resources/application.properties | 4 +- .../main/resources/mapper/FriendMapper.xml | 2 +- .../handlers/user/ReqAddFriendHandler.java | 23 ++++++--- .../handlers/user/ReqFriendListHandler.java | 26 +++++++++- .../soim/entry/worker/FutureExecutor.java | 50 +++++++++++++++++++ .../net/sopod/soim/entry/worker/Worker.java | 18 +++++++ .../route/ImRouterConsistentHashRoute.java | 1 + .../api/route/UidConsistentHashSelector.java | 14 +++++- .../logic/user/service/FriendServiceImpl.java | 23 +++++++++ 23 files changed, 279 insertions(+), 25 deletions(-) create mode 100644 im-client/src/main/java/net/sopod/soim/client/handler/cmd/FriendsHandler.java create mode 100644 im-client/src/main/java/net/sopod/soim/client/handler/msg/ResAddFriendHandler.java create mode 100644 im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendListHandler.java create mode 100644 im-entry/src/main/java/net/sopod/soim/entry/worker/FutureExecutor.java diff --git a/im-client/src/main/java/net/sopod/soim/client/cmd/CmdEnum.java b/im-client/src/main/java/net/sopod/soim/client/cmd/CmdEnum.java index fb03ee7..e165e8d 100644 --- a/im-client/src/main/java/net/sopod/soim/client/cmd/CmdEnum.java +++ b/im-client/src/main/java/net/sopod/soim/client/cmd/CmdEnum.java @@ -26,6 +26,11 @@ public enum CmdEnum { /** 用户搜索 */ search(SearchHandler.class), + /** 添加好友 */ + add(AddFriendHandler.class), + + friends(FriendsHandler.class), + help, /** 退出 */ diff --git a/im-client/src/main/java/net/sopod/soim/client/cmd/args/ArgsAddFriend.java b/im-client/src/main/java/net/sopod/soim/client/cmd/args/ArgsAddFriend.java index d4a7c4f..5b99c71 100644 --- a/im-client/src/main/java/net/sopod/soim/client/cmd/args/ArgsAddFriend.java +++ b/im-client/src/main/java/net/sopod/soim/client/cmd/args/ArgsAddFriend.java @@ -14,7 +14,7 @@ import java.util.List; @Data public class ArgsAddFriend { - @Parameter - private List parameters; + @Parameter(names = {"-u"}, required = true, description = "好友id") + private Long friendId; } diff --git a/im-client/src/main/java/net/sopod/soim/client/handler/cmd/AddFriendHandler.java b/im-client/src/main/java/net/sopod/soim/client/handler/cmd/AddFriendHandler.java index 7c7105b..84c5e70 100644 --- a/im-client/src/main/java/net/sopod/soim/client/handler/cmd/AddFriendHandler.java +++ b/im-client/src/main/java/net/sopod/soim/client/handler/cmd/AddFriendHandler.java @@ -2,7 +2,10 @@ package net.sopod.soim.client.handler.cmd; import com.google.inject.Inject; import com.google.inject.Singleton; +import net.sopod.soim.client.cmd.args.ArgsAddFriend; +import net.sopod.soim.client.cmd.handler.CmdHandler; import net.sopod.soim.client.session.SoImSession; +import net.sopod.soim.data.msg.user.Friend; /** * AddFriendHandler @@ -11,13 +14,22 @@ import net.sopod.soim.client.session.SoImSession; * @date 2022-05-21 19:29 */ @Singleton -public class AddFriendHandler { +public class AddFriendHandler implements CmdHandler { @Inject private SoImSession soImSession; - public void a() { + @Override + public ArgsAddFriend newArgsInstance() { + return new ArgsAddFriend(); + } + @Override + public void handleArgs(ArgsAddFriend args) { + Friend.ReqAddFriend req = Friend.ReqAddFriend.newBuilder() + .setFid(args.getFriendId()) + .build(); + soImSession.send(req); } } diff --git a/im-client/src/main/java/net/sopod/soim/client/handler/cmd/FriendsHandler.java b/im-client/src/main/java/net/sopod/soim/client/handler/cmd/FriendsHandler.java new file mode 100644 index 0000000..8f3d11b --- /dev/null +++ b/im-client/src/main/java/net/sopod/soim/client/handler/cmd/FriendsHandler.java @@ -0,0 +1,27 @@ +package net.sopod.soim.client.handler.cmd; + +import com.google.inject.Inject; +import com.google.inject.Singleton; +import net.sopod.soim.client.cmd.handler.NonArgsHandler; +import net.sopod.soim.client.session.SoImSession; +import net.sopod.soim.data.msg.user.Friend; + +/** + * FriendsHandler + * + * @author tmy + * @date 2022-05-23 22:58 + */ +@Singleton +public class FriendsHandler extends NonArgsHandler { + + @Inject + private SoImSession soImSession; + + @Override + public void handle() { + Friend.ReqFriendList req = Friend.ReqFriendList.newBuilder().build(); + soImSession.send(req); + } + +} diff --git a/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResAddFriendHandler.java b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResAddFriendHandler.java new file mode 100644 index 0000000..96a53f4 --- /dev/null +++ b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResAddFriendHandler.java @@ -0,0 +1,27 @@ +package net.sopod.soim.client.handler.msg; + +import com.google.inject.Singleton; +import net.sopod.soim.client.logger.Logger; +import net.sopod.soim.client.session.MessageHandler; +import net.sopod.soim.data.msg.user.Friend; + +/** + * ResAddFriendHandler + * + * @author tmy + * @date 2022-05-23 21:42 + */ +@Singleton +public class ResAddFriendHandler implements MessageHandler { + + @Override + public void handleMsg(Friend.ResAddFriend res) { + System.out.println(res); + if (res.getSuccess()) { + Logger.info("添加好友成功"); + } else { + Logger.error("添加失败:{}", res.getMsg()); + } + } + +} diff --git a/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendListHandler.java b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendListHandler.java new file mode 100644 index 0000000..5aede26 --- /dev/null +++ b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendListHandler.java @@ -0,0 +1,30 @@ +package net.sopod.soim.client.handler.msg; + +import com.google.inject.Singleton; +import net.sopod.soim.client.logger.Logger; +import net.sopod.soim.client.session.MessageHandler; +import net.sopod.soim.data.msg.user.Friend; +import net.sopod.soim.data.msg.user.UserMsg; + +import java.util.List; +import java.util.stream.Collectors; + +/** + * ResFriendListHandler + * + * @author tmy + * @date 2022-05-23 22:59 + */ +@Singleton +public class ResFriendListHandler implements MessageHandler { + + @Override + public void handleMsg(Friend.ResFriendList res) { + List friendsList = res.getFriendsList(); + List userLines = friendsList.stream() + .map(u -> u.getUid() + "|" + u.getAccount() + "|" + u.getNickname()) + .collect(Collectors.toList()); + Logger.logList("好友列表", userLines); + } + +} diff --git a/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendSearchHandler.java b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendSearchHandler.java index f4aa449..4ff49ff 100644 --- a/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendSearchHandler.java +++ b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendSearchHandler.java @@ -20,8 +20,8 @@ import java.util.stream.Collectors; public class ResFriendSearchHandler implements MessageHandler { @Override - public void handleMsg(AccountSearch.ResAccountSearch msg) { - List usersList = msg.getUsersList(); + public void handleMsg(AccountSearch.ResAccountSearch res) { + List usersList = res.getUsersList(); List userLines = usersList.stream() .map(u -> u.getUid() + "|" + u.getAccount() + "|" + u.getNickname()) .collect(Collectors.toList()); diff --git a/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResOnlineUserListHandler.java b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResOnlineUserListHandler.java index 1227852..a5d1101 100644 --- a/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResOnlineUserListHandler.java +++ b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResOnlineUserListHandler.java @@ -18,8 +18,8 @@ import java.util.stream.Collectors; public class ResOnlineUserListHandler implements MessageHandler { @Override - public void handleMsg(UserGroup.ResOnlineUserList msg) { - List userLines = msg.getUsersList().stream() + public void handleMsg(UserGroup.ResOnlineUserList res) { + List userLines = res.getUsersList().stream() .map(user -> user.getUid() + " | " + user.getAccount()) .collect(Collectors.toList()); Logger.logList("在线用户", userLines); diff --git a/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResTokenAuthHandler.java b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResTokenAuthHandler.java index 7e4334b..58b6d29 100644 --- a/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResTokenAuthHandler.java +++ b/im-client/src/main/java/net/sopod/soim/client/handler/msg/ResTokenAuthHandler.java @@ -19,8 +19,8 @@ public class ResTokenAuthHandler implements MessageHandler { private SoImSession soImSession; @Override - public void handleMsg(Auth.ResTokenAuth msg) { - soImSession.authResult(msg.getSuccess(), msg.getMessage(), msg.getUid()); + public void handleMsg(Auth.ResTokenAuth res) { + soImSession.authResult(res.getSuccess(), res.getMessage(), res.getUid()); } } diff --git a/im-client/src/main/java/net/sopod/soim/client/handler/msg/TextChatHandler.java b/im-client/src/main/java/net/sopod/soim/client/handler/msg/TextChatHandler.java index 6dc989d..2e0bda7 100644 --- a/im-client/src/main/java/net/sopod/soim/client/handler/msg/TextChatHandler.java +++ b/im-client/src/main/java/net/sopod/soim/client/handler/msg/TextChatHandler.java @@ -15,8 +15,8 @@ import net.sopod.soim.data.msg.chat.Chat; public class TextChatHandler implements MessageHandler { @Override - public void handleMsg(Chat.TextChat msg) { - Logger.info("{}: {}", msg.getSender(), msg.getMessage()); + public void handleMsg(Chat.TextChat res) { + Logger.info("{}: {}", res.getSender(), res.getMessage()); } } diff --git a/im-client/src/main/java/net/sopod/soim/client/session/MessageHandler.java b/im-client/src/main/java/net/sopod/soim/client/session/MessageHandler.java index 8d3d7f0..ff20f4e 100644 --- a/im-client/src/main/java/net/sopod/soim/client/session/MessageHandler.java +++ b/im-client/src/main/java/net/sopod/soim/client/session/MessageHandler.java @@ -8,6 +8,6 @@ package net.sopod.soim.client.session; */ public interface MessageHandler { - void handleMsg(T msg); + void handleMsg(T res); } diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/LogicTables.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/LogicTables.java index 3d01e82..0b30c67 100644 --- a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/LogicTables.java +++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/LogicTables.java @@ -10,4 +10,6 @@ public interface LogicTables { String IM_USER = "im_user"; + String IM_FRIEND = "im_friend"; + } diff --git a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/model/entity/ImFriend.java b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/model/entity/ImFriend.java index 014b96e..e717c64 100644 --- a/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/model/entity/ImFriend.java +++ b/im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/model/entity/ImFriend.java @@ -1,5 +1,6 @@ package net.sopod.soim.das.user.api.model.entity; +import com.baomidou.mybatisplus.annotation.IdType; import com.baomidou.mybatisplus.annotation.TableField; import com.baomidou.mybatisplus.annotation.TableId; import com.baomidou.mybatisplus.annotation.TableName; @@ -24,7 +25,7 @@ public class ImFriend implements Serializable { private static final long serialVersionUID = -4703347353944742873L; /** */ - @TableId(value = "id") + @TableId(value = "id", type = IdType.INPUT) private Long id; /** 用户id */ diff --git a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java index f37a685..6392518 100644 --- a/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java +++ b/im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java @@ -6,15 +6,18 @@ import lombok.AllArgsConstructor; import net.sopod.soim.common.constant.LogicConsts; import net.sopod.soim.common.util.Collects; import net.sopod.soim.common.util.ImClock; +import net.sopod.soim.das.user.api.config.LogicTables; import net.sopod.soim.das.user.api.model.entity.ImFriend; import net.sopod.soim.das.user.api.model.entity.ImUser; import net.sopod.soim.das.user.api.service.FriendDas; import net.sopod.soim.das.user.dao.FriendMapper; import net.sopod.soim.das.user.dao.UserMapper; +import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator; import org.apache.dubbo.config.annotation.DubboService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.Collections; import java.util.List; /** @@ -29,13 +32,18 @@ public class FriendDasImpl implements FriendDas { private static Logger logger = LoggerFactory.getLogger(FriendDasImpl.class); + SegmentIdGenerator segmentIdGenerator; + private FriendMapper friendMapper; private UserMapper userMapper; + @Override public int insert(Long uid, Long fid) { + long id = segmentIdGenerator.nextId(LogicTables.IM_FRIEND); ImFriend imFriend = new ImFriend() + .setId(id) .setUid(uid) .setFid(fid) .setStatus(1) @@ -71,6 +79,9 @@ public class FriendDasImpl implements FriendDas { @Override public List queryAllFriend(Long uid) { List fids = friendMapper.selectAllFriendId(uid); + if (Collects.isEmpty(fids)) { + return Collections.emptyList(); + } LambdaQueryWrapper userQuery = new QueryWrapper().lambda() .in(ImUser::getId, fids) .eq(ImUser::getStatus, LogicConsts.STATUS_NORMAL); diff --git a/im-das/im-das-user/src/main/resources/application.properties b/im-das/im-das-user/src/main/resources/application.properties index 63a7e3d..7382d45 100644 --- a/im-das/im-das-user/src/main/resources/application.properties +++ b/im-das/im-das-user/src/main/resources/application.properties @@ -25,11 +25,11 @@ spring.shardingsphere.rules.sharding.tables.im_user.table-strategy.standard.shar spring.shardingsphere.rules.sharding.sharding-algorithms.im-user-inline.type=INLINE spring.shardingsphere.rules.sharding.sharding-algorithms.im-user-inline.props.algorithm-expression=im_user_$->{id % 4} # im_friend -spring.shardingsphere.rules.sharding.tables.im_friend.actual-data-nodes=ds1.im_user_$->{0..7} +spring.shardingsphere.rules.sharding.tables.im_friend.actual-data-nodes=ds1.im_friend_$->{0..7} spring.shardingsphere.rules.sharding.tables.im_friend.table-strategy.standard.sharding-column=uid spring.shardingsphere.rules.sharding.tables.im_friend.table-strategy.standard.sharding-algorithm-name=im-friend-inline spring.shardingsphere.rules.sharding.sharding-algorithms.im-friend-inline.type=INLINE -spring.shardingsphere.rules.sharding.sharding-algorithms.im-friend-inline.props.algorithm-expression=im-friend-_$->{uid % 8} +spring.shardingsphere.rules.sharding.sharding-algorithms.im-friend-inline.props.algorithm-expression=im_friend_$->{uid % 8} # sharding-jdbc sql־ spring.shardingsphere.props.sql-show=true diff --git a/im-das/im-das-user/src/main/resources/mapper/FriendMapper.xml b/im-das/im-das-user/src/main/resources/mapper/FriendMapper.xml index a15823d..ba4f023 100644 --- a/im-das/im-das-user/src/main/resources/mapper/FriendMapper.xml +++ b/im-das/im-das-user/src/main/resources/mapper/FriendMapper.xml @@ -3,7 +3,7 @@ diff --git a/im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqAddFriendHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqAddFriendHandler.java index b605edc..fbbd661 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqAddFriendHandler.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqAddFriendHandler.java @@ -1,10 +1,12 @@ package net.sopod.soim.entry.handlers.user; import com.google.protobuf.MessageLite; -import net.sopod.soim.data.msg.user.AccountSearch; +import net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory; import net.sopod.soim.data.msg.user.Friend; import net.sopod.soim.entry.server.handler.AccountMessageHandler; import net.sopod.soim.entry.server.session.Account; +import net.sopod.soim.entry.worker.FutureExecutor; +import net.sopod.soim.entry.worker.WorkerGroup; import net.sopod.soim.logic.api.user.service.FriendService; import org.apache.dubbo.config.annotation.DubboReference; import org.slf4j.Logger; @@ -12,6 +14,8 @@ import org.slf4j.LoggerFactory; import org.springframework.stereotype.Service; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; /** * ReqAddFriendHandler @@ -34,16 +38,23 @@ public class ReqAddFriendHandler extends AccountMessageHandler addFuture = friendService.addFriend(uid, fid); - addFuture.whenComplete((msg, err) -> { + addFuture.whenCompleteAsync((msg, err) -> { if (err != null) { logger.error("添加好友失败:", err); return; } - Friend.ResAddFriend res = Friend.ResAddFriend.newBuilder() - .setSuccess(msg == null).setMsg(msg).build(); - account.writeNow(res); + Friend.ResAddFriend.Builder builder = Friend.ResAddFriend.newBuilder(); + builder.setSuccess(msg == null); + if (msg != null) { + builder.setMsg(msg); + } + account.writeNow(builder.build()); logger.info("添加好友成功"); - }); + }, FutureExecutor.getInstance()) + .exceptionally(err -> { + logger.error("future执行失败: ", err); + return null; + }); return null; } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqFriendListHandler.java b/im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqFriendListHandler.java index 451a9c7..06d1503 100644 --- a/im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqFriendListHandler.java +++ b/im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqFriendListHandler.java @@ -1,19 +1,29 @@ package net.sopod.soim.entry.handlers.user; import com.google.protobuf.MessageLite; +import lombok.extern.slf4j.Slf4j; import net.sopod.soim.data.msg.user.Friend; +import net.sopod.soim.data.msg.user.UserMsg; import net.sopod.soim.entry.server.handler.AccountMessageHandler; import net.sopod.soim.entry.server.session.Account; +import net.sopod.soim.entry.worker.FutureExecutor; import net.sopod.soim.logic.api.user.service.FriendService; +import net.sopod.soim.logic.common.model.UserInfo; import org.apache.dubbo.config.annotation.DubboReference; import org.springframework.stereotype.Service; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.stream.Collectors; + /** * ReqFriendListHandler * * @author tmy * @date 2022-05-23 17:42 */ +@Slf4j @Service public class ReqFriendListHandler extends AccountMessageHandler { @@ -22,7 +32,21 @@ public class ReqFriendListHandler extends AccountMessageHandler> future = friendService.listFriend(account.getUid()); + future.whenCompleteAsync((friends, err) -> { + List userInfos = friends.stream().map(acc -> UserMsg.UserInfo.newBuilder() + .setAccount(acc.getAccount()) + .setNickname(acc.getNickname()) + .setUid(acc.getUid()) + .build() + ).collect(Collectors.toList()); + Friend.ResFriendList res = Friend.ResFriendList.newBuilder().addAllFriends(userInfos).build(); + account.writeNow(res); + }, FutureExecutor.getInstance()) + .exceptionally(err -> { + log.error("好友列表查询失败", err); + return Collections.emptyList(); + }); return null; } diff --git a/im-entry/src/main/java/net/sopod/soim/entry/worker/FutureExecutor.java b/im-entry/src/main/java/net/sopod/soim/entry/worker/FutureExecutor.java new file mode 100644 index 0000000..02b729b --- /dev/null +++ b/im-entry/src/main/java/net/sopod/soim/entry/worker/FutureExecutor.java @@ -0,0 +1,50 @@ +package net.sopod.soim.entry.worker; + +import net.sopod.soim.common.util.netty.FastThreadLocalThreadFactory; + +import java.util.concurrent.Executor; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +/** + * FutureExecutor + * + * @author tmy + * @date 2022-05-23 23:03 + */ +public class FutureExecutor implements Executor { + + private static FutureExecutor INSTANCE; + + private final ThreadPoolExecutor threadPoolExecutor; + + private FutureExecutor() { + int cpus = Runtime.getRuntime().availableProcessors(); + threadPoolExecutor = new ThreadPoolExecutor( + 2, + cpus * 4, + 0, + TimeUnit.MILLISECONDS, + new LinkedBlockingQueue(), + new FastThreadLocalThreadFactory("future-exec-%d", Thread.NORM_PRIORITY) + ); + } + + @Override + public void execute(Runnable command) { + threadPoolExecutor.execute(command); + } + + public static FutureExecutor getInstance() { + if (INSTANCE == null) { + synchronized (FutureExecutor.class) { + if (INSTANCE == null) { + INSTANCE = new FutureExecutor(); + } + } + } + return INSTANCE; + } + +} 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 338177a..48fa9f9 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 @@ -9,6 +9,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.concurrent.*; +import java.util.function.BiConsumer; /** * dispatch -> worker * core_num -> disruptor 队列执行 @@ -76,4 +77,21 @@ public class Worker implements EventHandler, EventFactory public TaskEvent newInstance() { return new TaskEvent(); } + + public void whenCompleteAsync(CompletableFuture future, + BiConsumer action) { +// BiConsumer wrapper = (data, err) -> { +// try { +// action.accept(data, err); +// }catch (Exception e) { +// logger.error("future执行失败: ", e); +// } +// }; + future.whenCompleteAsync(action, executor) + .exceptionally(err -> { + logger.error("future执行失败: ", err); + return null; + }); + } + } diff --git a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ImRouterConsistentHashRoute.java b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ImRouterConsistentHashRoute.java index 18c1b90..7b01422 100644 --- a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ImRouterConsistentHashRoute.java +++ b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ImRouterConsistentHashRoute.java @@ -43,6 +43,7 @@ public class ImRouterConsistentHashRoute extends AbstractLoadBalance { invoker -> invoker.getUrl().getAddress(), new HashMap<>(6)); selector = new UidConsistentHashSelector<>(serverAddrInvokerMap, invokersHash); + selector.asCurrent(); } // 获取上下文 uid, 同 RpcContext String uid = invocation.getAttachment(DubboConstant.CTX_UID); diff --git a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java index be4ce17..da2c393 100644 --- a/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java +++ b/im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java @@ -1,6 +1,5 @@ package net.sopod.soim.router.api.route; -import net.sopod.soim.common.util.HashAlgorithms; import org.apache.commons.lang3.tuple.ImmutablePair; import org.apache.commons.lang3.tuple.Pair; @@ -120,4 +119,17 @@ public class UidConsistentHashSelector { return identityHashCode; } + private static volatile UidConsistentHashSelector currentSelector; + + void asCurrent() { + currentSelector = this; + } + + public static UidConsistentHashSelector getCurrent() { +// if (currentSelector == null) { +// throw new IllegalStateException("当前无运行中的UidConsistentHashSelector"); +// } + return currentSelector; + } + } diff --git a/im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/FriendServiceImpl.java b/im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/FriendServiceImpl.java index 858499c..d8fc5f0 100644 --- a/im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/FriendServiceImpl.java +++ b/im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/FriendServiceImpl.java @@ -1,14 +1,19 @@ package net.sopod.soim.logic.user.service; +import net.sopod.soim.common.util.Collects; +import net.sopod.soim.common.util.StringUtil; import net.sopod.soim.das.user.api.model.entity.ImUser; import net.sopod.soim.das.user.api.service.FriendDas; import net.sopod.soim.das.user.api.service.UserDas; import net.sopod.soim.logic.api.user.service.FriendService; import net.sopod.soim.logic.common.model.UserInfo; +import net.sopod.soim.logic.common.util.RpcContextUtil; +import net.sopod.soim.router.api.route.UidConsistentHashSelector; import org.apache.dubbo.config.annotation.DubboReference; import org.apache.dubbo.config.annotation.DubboService; import java.util.List; +import java.util.Map; import java.util.concurrent.CompletableFuture; import java.util.stream.Collectors; @@ -42,12 +47,30 @@ public class FriendServiceImpl implements FriendService { } // 添加好友数据 friendDas.insert(uid, fid); + // 相互添加为好友 + if (Boolean.FALSE.equals(friendDas.isExists(fid, uid))) { + friendDas.insert(fid, uid); + } return CompletableFuture.completedFuture(null); } @Override public CompletableFuture> listFriend(Long uid) { List imUsers = friendDas.queryAllFriend(uid); + // TODO 分组到 im-router 查询在线状态 + UidConsistentHashSelector selector = UidConsistentHashSelector.getCurrent(); + if (selector == null) { + // TODO 单实例 im-router, 直接调用 + } else { + Map> userRouteGroup = + Collects.group(imUsers, imUser -> selector.select(StringUtil.toString(imUser.getId())), imUser -> imUser); + for (List users : userRouteGroup.values()) { + List userIds = users.stream().map(ImUser::getId).collect(Collectors.toList()); + RpcContextUtil.setContextUid(userIds.get(0)); + // TODO 调用在线状态查询接口,按返回list索引更新在线状态 + + } + } List userInfos = imUsers.stream().map(imUser -> new UserInfo() .setUid(imUser.getId()) .setAccount(imUser.getAccount())