Browse Source

群聊消息转发

master
tangmingyou 4 years ago
parent
commit
53e2414fbd
  1. 4
      im-service/im-logic-message/pom.xml
  2. 37
      im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImGroupChatServiceImpl.java
  3. 23
      im-service/im-router/src/main/java/net/sopod/soim/router/service/MessageRouteServiceImpl.java
  4. 37
      im-service/im-router/src/main/java/net/sopod/soim/router/service/RouterUserService.java
  5. 1
      im-service/im-router/src/main/java/net/sopod/soim/router/util/RpcContextUtil.java

4
im-service/im-logic-message/pom.xml

@ -74,6 +74,10 @@
<groupId>org.apache.dubbo</groupId> <groupId>org.apache.dubbo</groupId>
<artifactId>dubbo-registry-nacos</artifactId> <artifactId>dubbo-registry-nacos</artifactId>
</dependency> </dependency>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
<dependency> <dependency>
<groupId>org.springframework.boot</groupId> <groupId>org.springframework.boot</groupId>

37
im-service/im-logic-message/src/main/java/net/sopod/soim/logic/message/service/ImGroupChatServiceImpl.java

@ -1,6 +1,8 @@
package net.sopod.soim.logic.message.service; package net.sopod.soim.logic.message.service;
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.common.util.StringUtil;
import net.sopod.soim.das.common.config.LogicTables; import net.sopod.soim.das.common.config.LogicTables;
import net.sopod.soim.das.group.api.model.dto.GroupUser_0; import net.sopod.soim.das.group.api.model.dto.GroupUser_0;
import net.sopod.soim.das.group.api.service.DasGroupUserService; import net.sopod.soim.das.group.api.service.DasGroupUserService;
@ -9,13 +11,19 @@ import net.sopod.soim.das.message.api.service.DasMQMessagePersistentService;
import net.sopod.soim.logic.api.message.mode.GroupMessage; import net.sopod.soim.logic.api.message.mode.GroupMessage;
import net.sopod.soim.logic.api.message.service.ImGroupChatService; import net.sopod.soim.logic.api.message.service.ImGroupChatService;
import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator; import net.sopod.soim.logic.api.segmentid.core.SegmentIdGenerator;
import net.sopod.soim.logic.common.util.RpcContextUtil;
import net.sopod.soim.router.api.route.UidConsistentHashSelector; import net.sopod.soim.router.api.route.UidConsistentHashSelector;
import net.sopod.soim.router.api.service.MessageRouteService;
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 javax.annotation.Resource; import javax.annotation.Resource;
import java.util.Iterator;
import java.util.List; import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletableFuture;
import java.util.stream.Collectors;
/** /**
* ImGroupChatServiceImpl * ImGroupChatServiceImpl
@ -35,6 +43,9 @@ public class ImGroupChatServiceImpl implements ImGroupChatService {
@DubboReference @DubboReference
private DasGroupUserService dasGroupUserService; private DasGroupUserService dasGroupUserService;
@DubboReference
private MessageRouteService messageRouteService;
/** /**
* 序列化群消息 * 序列化群消息
* 查询群成员列表 * 查询群成员列表
@ -54,11 +65,33 @@ public class ImGroupChatServiceImpl implements ImGroupChatService {
dasMQMessagePersistentService.saveImGroupMessage(imGroupMessage); dasMQMessagePersistentService.saveImGroupMessage(imGroupMessage);
// 查询群成员列表 // 查询群成员列表
List<GroupUser_0> groupUsers = dasGroupUserService.listGroupUsers(msg.getGid()); List<GroupUser_0> groupUsers = dasGroupUserService.listGroupUsers(msg.getGid());
Map<Long, GroupUser_0> groupUserMap = Collects.collect2Map(groupUsers, GroupUser_0::getUid);
// 分批路由消息 // 分批路由消息
UidConsistentHashSelector<?> uidSelector = UidConsistentHashSelector.getCurrent(); UidConsistentHashSelector<?> uidSelector = UidConsistentHashSelector.getCurrent();
if (uidSelector == null) {
List<Long> uidList = groupUsers.stream().map(GroupUser_0::getUid).collect(Collectors.toList());
List<Boolean> results = messageRouteService.routeGroupMessage(uidList, msg.getMessage());
// TODO 更新未读消息数据偏移量,未读消息条数
return CompletableFuture.completedFuture("OK");
}
Map<?, List<Long>> uidGroups = Collects.group(
groupUsers,
user -> uidSelector.select(StringUtil.toString(user.getUid())),
GroupUser_0::getUid
);
// 批量路由消息到对应的 router
for (List<Long> uidGroup : uidGroups.values()) {
RpcContextUtil.setContextUid(uidGroup.get(0));
List<Boolean> results = messageRouteService.routeGroupMessage(uidGroup, msg.getMessage());
Iterator<Long> iterator = uidGroup.iterator();
for (int i = 0; iterator.hasNext(); i++) {
GroupUser_0 gUser = groupUserMap.get(iterator.next());
Boolean received = results.get(i);
// TODO 更新未读消息数据偏移量,未读消息条数
}
return null; }
return CompletableFuture.completedFuture("OK");
} }
} }

23
im-service/im-router/src/main/java/net/sopod/soim/router/service/MessageRouteServiceImpl.java

@ -1,30 +1,43 @@
package net.sopod.soim.router.service; package net.sopod.soim.router.service;
import lombok.AllArgsConstructor;
import net.sopod.soim.router.api.service.MessageRouteService; import net.sopod.soim.router.api.service.MessageRouteService;
import net.sopod.soim.router.cache.RouterUser; import net.sopod.soim.router.cache.RouterUser;
import net.sopod.soim.router.cache.RouterUserStorage; import net.sopod.soim.router.cache.RouterUserStorage;
import org.apache.dubbo.config.annotation.DubboService; import org.apache.dubbo.config.annotation.DubboService;
import java.util.ArrayList;
import java.util.Iterator;
import java.util.List; import java.util.List;
/** /**
* MessageRouterServiceImpl * MessageRouterServiceImpl
* 路由聊天消息
* *
* @author tmy * @author tmy
* @date 2022-06-05 11:30 * @date 2022-06-05 11:30
*/ */
@DubboService @DubboService
@AllArgsConstructor
public class MessageRouteServiceImpl implements MessageRouteService { public class MessageRouteServiceImpl implements MessageRouteService {
private RouterUserService routerUserService;
@Override @Override
public List<Boolean> routeGroupMessage(List<Long> uids, String message) { public List<Boolean> routeGroupMessage(List<Long> uids, String message) {
// TODO
RouterUserStorage storage = RouterUserStorage.getInstance(); RouterUserStorage storage = RouterUserStorage.getInstance();
for (Long uid : uids) { List<Boolean> results = new ArrayList<>(uids.size());
RouterUser routerUser = storage.get(uid); Iterator<Long> iterator = uids.iterator();
for (int i = 0; iterator.hasNext(); i++) {
RouterUser routerUser = storage.get(iterator.next());
if (routerUser == null) {
results.set(i, false);
continue;
}
Boolean success = routerUserService.routeGroupMessage(routerUser, message);
results.set(i, Boolean.TRUE.equals(success));
} }
return null; return results;
} }
} }

37
im-service/im-router/src/main/java/net/sopod/soim/router/service/RouterUserService.java

@ -0,0 +1,37 @@
package net.sopod.soim.router.service;
import net.sopod.soim.common.util.ImClock;
import net.sopod.soim.entry.api.service.TextChatService;
import net.sopod.soim.logic.common.model.TextChat;
import net.sopod.soim.logic.common.util.RpcContextUtil;
import net.sopod.soim.router.cache.RouterUser;
import org.apache.dubbo.config.annotation.DubboReference;
import org.springframework.stereotype.Service;
/**
* RouterUserService
*
* @author tmy
* @date 2022-06-06 9:34
*/
@Service
public class RouterUserService {
@DubboReference
private TextChatService textChatService;
public void routeChatMessage(RouterUser routerUser, String message) {
}
public Boolean routeGroupMessage(RouterUser receiverUser, String message) {
RpcContextUtil.setContextUid(receiverUser.getUid());
TextChat textChat = new TextChat()
.setReceiverUid(receiverUser.getUid())
.setMessage(message)
.setTime(ImClock.millis())
.setReceiverName(receiverUser.getAccount());
return textChatService.sendTextChat(textChat);
}
}

1
im-service/im-router/src/main/java/net/sopod/soim/router/util/RpcContextUtil.java

@ -13,6 +13,7 @@ import java.util.Objects;
* @author tmy * @author tmy
* @date 2022-05-02 22:27 * @date 2022-05-02 22:27
*/ */
@Deprecated
public class RpcContextUtil { public class RpcContextUtil {
public static boolean setImEntryRouteServerAddrByUid(Long uid) { public static boolean setImEntryRouteServerAddrByUid(Long uid) {

Loading…
Cancel
Save