Browse Source

添加好友,好友列表

master
tangmingyou 4 years ago
parent
commit
17145b7ba2
  1. 5
      im-client/src/main/java/net/sopod/soim/client/cmd/CmdEnum.java
  2. 4
      im-client/src/main/java/net/sopod/soim/client/cmd/args/ArgsAddFriend.java
  3. 16
      im-client/src/main/java/net/sopod/soim/client/handler/cmd/AddFriendHandler.java
  4. 27
      im-client/src/main/java/net/sopod/soim/client/handler/cmd/FriendsHandler.java
  5. 27
      im-client/src/main/java/net/sopod/soim/client/handler/msg/ResAddFriendHandler.java
  6. 30
      im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendListHandler.java
  7. 4
      im-client/src/main/java/net/sopod/soim/client/handler/msg/ResFriendSearchHandler.java
  8. 4
      im-client/src/main/java/net/sopod/soim/client/handler/msg/ResOnlineUserListHandler.java
  9. 4
      im-client/src/main/java/net/sopod/soim/client/handler/msg/ResTokenAuthHandler.java
  10. 4
      im-client/src/main/java/net/sopod/soim/client/handler/msg/TextChatHandler.java
  11. 2
      im-client/src/main/java/net/sopod/soim/client/session/MessageHandler.java
  12. 2
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/config/LogicTables.java
  13. 3
      im-das-api/im-das-user-api/src/main/java/net/sopod/soim/das/user/api/model/entity/ImFriend.java
  14. 11
      im-das/im-das-user/src/main/java/net/sopod/soim/das/user/service/FriendDasImpl.java
  15. 4
      im-das/im-das-user/src/main/resources/application.properties
  16. 2
      im-das/im-das-user/src/main/resources/mapper/FriendMapper.xml
  17. 23
      im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqAddFriendHandler.java
  18. 26
      im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqFriendListHandler.java
  19. 50
      im-entry/src/main/java/net/sopod/soim/entry/worker/FutureExecutor.java
  20. 18
      im-entry/src/main/java/net/sopod/soim/entry/worker/Worker.java
  21. 1
      im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/ImRouterConsistentHashRoute.java
  22. 14
      im-service-api/im-router-api/src/main/java/net/sopod/soim/router/api/route/UidConsistentHashSelector.java
  23. 23
      im-service/im-logic-user/src/main/java/net/sopod/soim/logic/user/service/FriendServiceImpl.java

5
im-client/src/main/java/net/sopod/soim/client/cmd/CmdEnum.java

@ -26,6 +26,11 @@ public enum CmdEnum {
/** 用户搜索 */ /** 用户搜索 */
search(SearchHandler.class), search(SearchHandler.class),
/** 添加好友 */
add(AddFriendHandler.class),
friends(FriendsHandler.class),
help, help,
/** 退出 */ /** 退出 */

4
im-client/src/main/java/net/sopod/soim/client/cmd/args/ArgsAddFriend.java

@ -14,7 +14,7 @@ import java.util.List;
@Data @Data
public class ArgsAddFriend { public class ArgsAddFriend {
@Parameter @Parameter(names = {"-u"}, required = true, description = "好友id")
private List<String> parameters; private Long friendId;
} }

16
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.Inject;
import com.google.inject.Singleton; 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.client.session.SoImSession;
import net.sopod.soim.data.msg.user.Friend;
/** /**
* AddFriendHandler * AddFriendHandler
@ -11,13 +14,22 @@ import net.sopod.soim.client.session.SoImSession;
* @date 2022-05-21 19:29 * @date 2022-05-21 19:29
*/ */
@Singleton @Singleton
public class AddFriendHandler { public class AddFriendHandler implements CmdHandler<ArgsAddFriend> {
@Inject @Inject
private SoImSession soImSession; 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);
} }
} }

27
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);
}
}

27
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<Friend.ResAddFriend> {
@Override
public void handleMsg(Friend.ResAddFriend res) {
System.out.println(res);
if (res.getSuccess()) {
Logger.info("添加好友成功");
} else {
Logger.error("添加失败:{}", res.getMsg());
}
}
}

30
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<Friend.ResFriendList> {
@Override
public void handleMsg(Friend.ResFriendList res) {
List<UserMsg.UserInfo> friendsList = res.getFriendsList();
List<String> userLines = friendsList.stream()
.map(u -> u.getUid() + "|" + u.getAccount() + "|" + u.getNickname())
.collect(Collectors.toList());
Logger.logList("好友列表", userLines);
}
}

4
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<AccountSearch.ResAccountSearch> { public class ResFriendSearchHandler implements MessageHandler<AccountSearch.ResAccountSearch> {
@Override @Override
public void handleMsg(AccountSearch.ResAccountSearch msg) { public void handleMsg(AccountSearch.ResAccountSearch res) {
List<UserMsg.UserInfo> usersList = msg.getUsersList(); List<UserMsg.UserInfo> usersList = res.getUsersList();
List<String> userLines = usersList.stream() List<String> userLines = usersList.stream()
.map(u -> u.getUid() + "|" + u.getAccount() + "|" + u.getNickname()) .map(u -> u.getUid() + "|" + u.getAccount() + "|" + u.getNickname())
.collect(Collectors.toList()); .collect(Collectors.toList());

4
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<UserGroup.ResOnlineUserList> { public class ResOnlineUserListHandler implements MessageHandler<UserGroup.ResOnlineUserList> {
@Override @Override
public void handleMsg(UserGroup.ResOnlineUserList msg) { public void handleMsg(UserGroup.ResOnlineUserList res) {
List<String> userLines = msg.getUsersList().stream() List<String> userLines = res.getUsersList().stream()
.map(user -> user.getUid() + " | " + user.getAccount()) .map(user -> user.getUid() + " | " + user.getAccount())
.collect(Collectors.toList()); .collect(Collectors.toList());
Logger.logList("在线用户", userLines); Logger.logList("在线用户", userLines);

4
im-client/src/main/java/net/sopod/soim/client/handler/msg/ResTokenAuthHandler.java

@ -19,8 +19,8 @@ public class ResTokenAuthHandler implements MessageHandler<Auth.ResTokenAuth> {
private SoImSession soImSession; private SoImSession soImSession;
@Override @Override
public void handleMsg(Auth.ResTokenAuth msg) { public void handleMsg(Auth.ResTokenAuth res) {
soImSession.authResult(msg.getSuccess(), msg.getMessage(), msg.getUid()); soImSession.authResult(res.getSuccess(), res.getMessage(), res.getUid());
} }
} }

4
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<Chat.TextChat> { public class TextChatHandler implements MessageHandler<Chat.TextChat> {
@Override @Override
public void handleMsg(Chat.TextChat msg) { public void handleMsg(Chat.TextChat res) {
Logger.info("{}: {}", msg.getSender(), msg.getMessage()); Logger.info("{}: {}", res.getSender(), res.getMessage());
} }
} }

2
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<T> { public interface MessageHandler<T> {
void handleMsg(T msg); void handleMsg(T res);
} }

2
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_USER = "im_user";
String IM_FRIEND = "im_friend";
} }

3
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; 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.TableField;
import com.baomidou.mybatisplus.annotation.TableId; import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName; import com.baomidou.mybatisplus.annotation.TableName;
@ -24,7 +25,7 @@ public class ImFriend implements Serializable {
private static final long serialVersionUID = -4703347353944742873L; private static final long serialVersionUID = -4703347353944742873L;
/** */ /** */
@TableId(value = "id") @TableId(value = "id", type = IdType.INPUT)
private Long id; private Long id;
/** 用户id */ /** 用户id */

11
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.constant.LogicConsts;
import net.sopod.soim.common.util.Collects; import net.sopod.soim.common.util.Collects;
import net.sopod.soim.common.util.ImClock; 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.ImFriend;
import net.sopod.soim.das.user.api.model.entity.ImUser; 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.FriendDas;
import net.sopod.soim.das.user.dao.FriendMapper; import net.sopod.soim.das.user.dao.FriendMapper;
import net.sopod.soim.das.user.dao.UserMapper; 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.apache.dubbo.config.annotation.DubboService;
import org.slf4j.Logger; import org.slf4j.Logger;
import org.slf4j.LoggerFactory; import org.slf4j.LoggerFactory;
import java.util.Collections;
import java.util.List; import java.util.List;
/** /**
@ -29,13 +32,18 @@ public class FriendDasImpl implements FriendDas {
private static Logger logger = LoggerFactory.getLogger(FriendDasImpl.class); private static Logger logger = LoggerFactory.getLogger(FriendDasImpl.class);
SegmentIdGenerator segmentIdGenerator;
private FriendMapper friendMapper; private FriendMapper friendMapper;
private UserMapper userMapper; private UserMapper userMapper;
@Override @Override
public int insert(Long uid, Long fid) { public int insert(Long uid, Long fid) {
long id = segmentIdGenerator.nextId(LogicTables.IM_FRIEND);
ImFriend imFriend = new ImFriend() ImFriend imFriend = new ImFriend()
.setId(id)
.setUid(uid) .setUid(uid)
.setFid(fid) .setFid(fid)
.setStatus(1) .setStatus(1)
@ -71,6 +79,9 @@ public class FriendDasImpl implements FriendDas {
@Override @Override
public List<ImUser> queryAllFriend(Long uid) { public List<ImUser> queryAllFriend(Long uid) {
List<Long> fids = friendMapper.selectAllFriendId(uid); List<Long> fids = friendMapper.selectAllFriendId(uid);
if (Collects.isEmpty(fids)) {
return Collections.emptyList();
}
LambdaQueryWrapper<ImUser> userQuery = new QueryWrapper<ImUser>().lambda() LambdaQueryWrapper<ImUser> userQuery = new QueryWrapper<ImUser>().lambda()
.in(ImUser::getId, fids) .in(ImUser::getId, fids)
.eq(ImUser::getStatus, LogicConsts.STATUS_NORMAL); .eq(ImUser::getStatus, LogicConsts.STATUS_NORMAL);

4
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.type=INLINE
spring.shardingsphere.rules.sharding.sharding-algorithms.im-user-inline.props.algorithm-expression=im_user_$->{id % 4} spring.shardingsphere.rules.sharding.sharding-algorithms.im-user-inline.props.algorithm-expression=im_user_$->{id % 4}
# im_friend # 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-column=uid
spring.shardingsphere.rules.sharding.tables.im_friend.table-strategy.standard.sharding-algorithm-name=im-friend-inline 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.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Êä³öÈÕÖ¾ # ´ò¿ªsharding-jdbc sqlÊä³öÈÕÖ¾
spring.shardingsphere.props.sql-show=true spring.shardingsphere.props.sql-show=true

2
im-das/im-das-user/src/main/resources/mapper/FriendMapper.xml

@ -3,7 +3,7 @@
<mapper namespace="net.sopod.soim.das.user.dao.FriendMapper"> <mapper namespace="net.sopod.soim.das.user.dao.FriendMapper">
<select id="selectAllFriendId" resultType="long"> <select id="selectAllFriendId" resultType="long">
select fid from im_friend where uid = #{uid} and status = 0 select fid from im_friend where uid = #{uid} and status = 1
</select> </select>
</mapper> </mapper>

23
im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqAddFriendHandler.java

@ -1,10 +1,12 @@
package net.sopod.soim.entry.handlers.user; package net.sopod.soim.entry.handlers.user;
import com.google.protobuf.MessageLite; 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.data.msg.user.Friend;
import net.sopod.soim.entry.server.handler.AccountMessageHandler; import net.sopod.soim.entry.server.handler.AccountMessageHandler;
import net.sopod.soim.entry.server.session.Account; 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 net.sopod.soim.logic.api.user.service.FriendService;
import org.apache.dubbo.config.annotation.DubboReference; import org.apache.dubbo.config.annotation.DubboReference;
import org.slf4j.Logger; import org.slf4j.Logger;
@ -12,6 +14,8 @@ import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
/** /**
* ReqAddFriendHandler * ReqAddFriendHandler
@ -34,16 +38,23 @@ public class ReqAddFriendHandler extends AccountMessageHandler<Friend.ReqAddFrie
long uid = account.getUid(); long uid = account.getUid();
logger.info("添加好友: {}, {}", uid, fid); logger.info("添加好友: {}, {}", uid, fid);
CompletableFuture<String> addFuture = friendService.addFriend(uid, fid); CompletableFuture<String> addFuture = friendService.addFriend(uid, fid);
addFuture.whenComplete((msg, err) -> { addFuture.whenCompleteAsync((msg, err) -> {
if (err != null) { if (err != null) {
logger.error("添加好友失败:", err); logger.error("添加好友失败:", err);
return; return;
} }
Friend.ResAddFriend res = Friend.ResAddFriend.newBuilder() Friend.ResAddFriend.Builder builder = Friend.ResAddFriend.newBuilder();
.setSuccess(msg == null).setMsg(msg).build(); builder.setSuccess(msg == null);
account.writeNow(res); if (msg != null) {
builder.setMsg(msg);
}
account.writeNow(builder.build());
logger.info("添加好友成功"); logger.info("添加好友成功");
}); }, FutureExecutor.getInstance())
.exceptionally(err -> {
logger.error("future执行失败: ", err);
return null;
});
return null; return null;
} }

26
im-entry/src/main/java/net/sopod/soim/entry/handlers/user/ReqFriendListHandler.java

@ -1,19 +1,29 @@
package net.sopod.soim.entry.handlers.user; package net.sopod.soim.entry.handlers.user;
import com.google.protobuf.MessageLite; 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.Friend;
import net.sopod.soim.data.msg.user.UserMsg;
import net.sopod.soim.entry.server.handler.AccountMessageHandler; import net.sopod.soim.entry.server.handler.AccountMessageHandler;
import net.sopod.soim.entry.server.session.Account; 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.api.user.service.FriendService;
import net.sopod.soim.logic.common.model.UserInfo;
import org.apache.dubbo.config.annotation.DubboReference; import org.apache.dubbo.config.annotation.DubboReference;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.stream.Collectors;
/** /**
* ReqFriendListHandler * ReqFriendListHandler
* *
* @author tmy * @author tmy
* @date 2022-05-23 17:42 * @date 2022-05-23 17:42
*/ */
@Slf4j
@Service @Service
public class ReqFriendListHandler extends AccountMessageHandler<Friend.ReqFriendList> { public class ReqFriendListHandler extends AccountMessageHandler<Friend.ReqFriendList> {
@ -22,7 +32,21 @@ public class ReqFriendListHandler extends AccountMessageHandler<Friend.ReqFriend
@Override @Override
public MessageLite handle(Account account, Friend.ReqFriendList req) { public MessageLite handle(Account account, Friend.ReqFriendList req) {
CompletableFuture<List<UserInfo>> future = friendService.listFriend(account.getUid());
future.whenCompleteAsync((friends, err) -> {
List<UserMsg.UserInfo> 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; return null;
} }

50
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<Runnable>(),
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;
}
}

18
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 org.slf4j.LoggerFactory;
import java.util.concurrent.*; import java.util.concurrent.*;
import java.util.function.BiConsumer;
/** /**
* dispatch -> worker * core_num -> disruptor 队列执行 * dispatch -> worker * core_num -> disruptor 队列执行
@ -76,4 +77,21 @@ public class Worker implements EventHandler<TaskEvent>, EventFactory<TaskEvent>
public TaskEvent newInstance() { public TaskEvent newInstance() {
return new TaskEvent(); return new TaskEvent();
} }
public <T> void whenCompleteAsync(CompletableFuture<T> future,
BiConsumer<? super T, ? super Throwable> action) {
// BiConsumer<? super T, ? super Throwable> 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;
});
}
} }

1
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(), invoker -> invoker.getUrl().getAddress(),
new HashMap<>(6)); new HashMap<>(6));
selector = new UidConsistentHashSelector<>(serverAddrInvokerMap, invokersHash); selector = new UidConsistentHashSelector<>(serverAddrInvokerMap, invokersHash);
selector.asCurrent();
} }
// 获取上下文 uid, 同 RpcContext // 获取上下文 uid, 同 RpcContext
String uid = invocation.getAttachment(DubboConstant.CTX_UID); String uid = invocation.getAttachment(DubboConstant.CTX_UID);

14
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; 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.ImmutablePair;
import org.apache.commons.lang3.tuple.Pair; import org.apache.commons.lang3.tuple.Pair;
@ -120,4 +119,17 @@ public class UidConsistentHashSelector<V> {
return identityHashCode; return identityHashCode;
} }
private static volatile UidConsistentHashSelector<?> currentSelector;
void asCurrent() {
currentSelector = this;
}
public static UidConsistentHashSelector<?> getCurrent() {
// if (currentSelector == null) {
// throw new IllegalStateException("当前无运行中的UidConsistentHashSelector");
// }
return currentSelector;
}
} }

23
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; 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.model.entity.ImUser;
import net.sopod.soim.das.user.api.service.FriendDas; import net.sopod.soim.das.user.api.service.FriendDas;
import net.sopod.soim.das.user.api.service.UserDas; import net.sopod.soim.das.user.api.service.UserDas;
import net.sopod.soim.logic.api.user.service.FriendService; import net.sopod.soim.logic.api.user.service.FriendService;
import net.sopod.soim.logic.common.model.UserInfo; 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.DubboReference;
import org.apache.dubbo.config.annotation.DubboService; import org.apache.dubbo.config.annotation.DubboService;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.stream.Collectors; import java.util.stream.Collectors;
@ -42,12 +47,30 @@ public class FriendServiceImpl implements FriendService {
} }
// 添加好友数据 // 添加好友数据
friendDas.insert(uid, fid); friendDas.insert(uid, fid);
// 相互添加为好友
if (Boolean.FALSE.equals(friendDas.isExists(fid, uid))) {
friendDas.insert(fid, uid);
}
return CompletableFuture.completedFuture(null); return CompletableFuture.completedFuture(null);
} }
@Override @Override
public CompletableFuture<List<UserInfo>> listFriend(Long uid) { public CompletableFuture<List<UserInfo>> listFriend(Long uid) {
List<ImUser> imUsers = friendDas.queryAllFriend(uid); List<ImUser> imUsers = friendDas.queryAllFriend(uid);
// TODO 分组到 im-router 查询在线状态
UidConsistentHashSelector<?> selector = UidConsistentHashSelector.getCurrent();
if (selector == null) {
// TODO 单实例 im-router, 直接调用
} else {
Map<?, List<ImUser>> userRouteGroup =
Collects.group(imUsers, imUser -> selector.select(StringUtil.toString(imUser.getId())), imUser -> imUser);
for (List<ImUser> users : userRouteGroup.values()) {
List<Long> userIds = users.stream().map(ImUser::getId).collect(Collectors.toList());
RpcContextUtil.setContextUid(userIds.get(0));
// TODO 调用在线状态查询接口,按返回list索引更新在线状态
}
}
List<UserInfo> userInfos = imUsers.stream().map(imUser -> new UserInfo() List<UserInfo> userInfos = imUsers.stream().map(imUser -> new UserInfo()
.setUid(imUser.getId()) .setUid(imUser.getId())
.setAccount(imUser.getAccount()) .setAccount(imUser.getAccount())

Loading…
Cancel
Save