Skip to content
第 247 / 250 章架构⏱ 18 分钟阅读

第 247 章:综合项目实战 - 即时通讯平台

学习目标

  • 长连接架构设计
  • 消息可靠投递
  • 群聊与单聊实现
  • 离线推送

一、IM 系统架构

二、长连接实现

2.1 Netty 服务端

java
@Component
@Slf4j
public class NettyServer {

    @Autowired
    private ChannelManager channelManager;

    @Autowired
    private MessageHandler messageHandler;

    @PostConstruct
    public void start() throws Exception {
        EventLoopGroup boss = new NioEventLoopGroup(1);
        EventLoopGroup worker = new NioEventLoopGroup(8);

        try {
            ServerBootstrap b = new ServerBootstrap();
            b.group(boss, worker)
                .channel(NioServerSocketChannel.class)
                .childOption(ChannelOption.TCP_NODELAY, true)
                .childOption(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT)
                .childHandler(new ChannelInitializer<SocketChannel>() {
                    @Override
                    protected void initChannel(SocketChannel ch) {
                        ch.pipeline()
                            .addLast(new LengthFieldBasedFrameDecoder(1024 * 1024, 0, 4, 0, 4))
                            .addLast(new LengthFieldPrepender(4))
                            .addLast(new IdleStateHandler(60, 30, 0))
                            .addLast(new AuthHandler())
                            .addLast(new MessageDecoder())
                            .addLast(messageHandler);
                    }
                });

            ChannelFuture cf = b.bind(8081).sync();
            cf.channel().closeFuture().sync();
        } finally {
            boss.shutdownGracefully();
            worker.shutdownGracefully();
        }
    }
}

2.2 Channel 管理

java
@Component
public class ChannelManager {

    private final ConcurrentHashMap<Long, Channel> userChannels = new ConcurrentHashMap<>();

    public void register(Long userId, Channel channel) {
        userChannels.put(userId, channel);
        log.info("用户 {} 注册连接,当前连接数:{}", userId, userChannels.size());
    }

    public void unregister(Long userId) {
        userChannels.remove(userId);
    }

    public void sendToUser(Long userId, Message msg) {
        Channel channel = userChannels.get(userId);
        if (channel != null && channel.isActive()) {
            channel.writeAndFlush(msg);
        } else {
            // 用户离线,存离线消息
            offlineMessageService.save(userId, msg);
        }
    }

    public boolean isOnline(Long userId) {
        Channel channel = userChannels.get(userId);
        return channel != null && channel.isActive();
    }
}

2.3 消息协议

java
@Data
public class Message {
    /** 消息类型 */
    private MessageType type;

    /** 发送者 */
    private Long fromUserId;

    /** 接收者(单聊 = userId, 群聊 = groupId) */
    private Long toId;

    /** 客户端消息 ID(用于去重) */
    private String clientMsgId;

    /** 服务端消息 ID */
    private String serverMsgId;

    /** 会话类型 */
    private ConversationType conversationType;

    /** 内容 */
    private String content;

    /** 媒体 URL */
    private String mediaUrl;

    /** 时间戳 */
    private Long timestamp;
}

public enum MessageType {
    HEARTBEAT,
    AUTH,
    TEXT,
    IMAGE,
    VOICE,
    VIDEO,
    FILE,
    ACK,
    READ,
    RECALL,
    TYPING
}

三、消息可靠性

3.1 消息状态机

3.2 可靠投递

java
@Service
@Slf4j
public class ReliableMessageService {

    @Autowired
    private ChannelManager channelManager;

    @Autowired
    private MessageStoreService storeService;

    @Autowired
    private KafkaTemplate<String, Message> kafkaTemplate;

    /**
     * 发送消息(完整流程)
     */
    public SendResult send(Message message) {
        // 1. 生成服务端 ID
        message.setServerMsgId(UUID.randomUUID().toString());
        message.setTimestamp(System.currentTimeMillis());

        // 2. 持久化(消息状态 Sending)
        storeService.saveWithStatus(message, MessageStatus.SENDING);

        // 3. 异步投递
        CompletableFuture.runAsync(() -> {
            try {
                deliver(message);
            } catch (Exception e) {
                log.error("消息投递失败", e);
            }
        });

        return new SendResult(message.getServerMsgId());
    }

    private void deliver(Message message) {
        if (message.getConversationType() == ConversationType.SINGLE) {
            // 单聊:通过 channel push
            channelManager.sendToUser(message.getToId(), message);

        } else if (message.getConversationType() == ConversationType.GROUP) {
            // 群聊:取所有成员
            List<Long> members = groupService.getMembers(message.getToId());
            for (Long member : members) {
                if (member.equals(message.getFromUserId())) continue;
                channelManager.sendToUser(member, message);
            }
        }

        // 4. 更新状态
        storeService.updateStatus(message.getServerMsgId(), MessageStatus.SENT);
    }

    /**
     * ACK 处理
     */
    public void handleAck(String serverMsgId, Long userId) {
        storeService.updateDelivered(serverMsgId, userId);
    }

    /**
     * 已读处理
     */
    public void handleRead(String serverMsgId, Long userId) {
        storeService.updateRead(serverMsgId, userId);
    }
}

3.3 ACK 机制

java
@Component
@Slf4j
public class AckManager {

    /** 待 ACK 消息 */
    private final ConcurrentHashMap<String, PendingAck> pendingAcks = new ConcurrentHashMap<>();

    /** 三秒超时重发 */
    private static final long TIMEOUT_MS = 3000;

    public void addPending(Message msg) {
        PendingAck pending = new PendingAck(msg, TIMEOUT_MS);
        pendingAcks.put(msg.getServerMsgId(), pending);

        // 启动超时检查
        pending.schedule(this::onTimeout);
    }

    public void onAck(String serverMsgId) {
        PendingAck pending = pendingAcks.remove(serverMsgId);
        if (pending != null) {
            pending.cancel();
        }
    }

    private void onTimeout(PendingAck pending) {
        if (pending.getAttempts() >= 3) {
            log.warn("消息 {} 重试 {} 次仍失败", pending.getMessage().getServerMsgId(),
                    pending.getAttempts());
            return;
        }

        // 重发
        reliableMessageService.deliver(pending.getMessage());
        pendingAcks.put(pending.getMessage().getServerMsgId(), pending);
    }
}

四、单聊与群聊

4.1 消息存储

java
@Service
public class MessageStoreService {

    @Autowired
    private MessageMapper messageMapper;

    @Autowired
    private ShardingJdbcTemplate jdbcTemplate;

    /**
     * 存储消息(按 userId 分库)
     */
    public void saveWithStatus(Message msg, MessageStatus status) {
        MessageEntity entity = new MessageEntity();
        BeanUtils.copyProperties(msg, entity);
        entity.setStatus(status);
        entity.setCreatedAt(LocalDateTime.now());
        jdbcTemplate.insert(entity);
    }

    /**
     * 拉取消息
     */
    public List<Message> pull(Long userId, Long conversationId, Long fromId, int limit) {
        return jdbcTemplate.query(
            "SELECT * FROM t_message WHERE conversation_id = ? AND id > ? ORDER BY id DESC LIMIT ?",
            conversationId, fromId, limit
        );
    }
}

4.2 会话列表

java
@Service
public class ConversationService {

    @Autowired
    private RedisTemplate<String, String> redisTemplate;

    /**
     * 获取会话列表(按最后消息时间排序)
     */
    public List<Conversation> list(Long userId) {
        String key = "conv:" + userId;
        Set<TypedTuple<String>> result = redisTemplate.opsForZSet()
            .reverseRangeWithScores(key, 0, -1);

        if (result == null) return Collections.emptyList();

        return result.stream()
            .map(t -> {
                Conversation conv = JSON.parseObject(t.getValue(), Conversation.class);
                conv.setLastActiveAt(new Date(t.getScore().longValue()));
                return conv;
            })
            .collect(Collectors.toList());
    }

    /**
     * 新消息后更新会话
     */
    public void update(Message msg) {
        // 更新发送方
        String fromKey = "conv:" + msg.getFromUserId();
        updateConversation(fromKey, msg);

        // 更新接收方
        String toKey = "conv:" + msg.getToId();
        updateConversation(toKey, msg);
    }

    private void updateConversation(String key, Message msg) {
        Conversation conv = new Conversation();
        conv.setConversationId(msg.getConversationType() == ConversationType.SINGLE
            ? msg.getToId() : msg.getToId());
        conv.setType(msg.getConversationType());
        conv.setLastMessage(msg.getContent());
        conv.setUnread(redisTemplate.opsForValue().increment(key + ":unread") != null ?
            redisTemplate.opsForValue().increment(key + ":unread").intValue() : 1);

        redisTemplate.opsForZSet().add(key, JSON.toJSONString(conv),
            System.currentTimeMillis());

        // 已读清除
        conv.setUnread(0);
    }
}

4.3 群聊实现

java
@Service
public class GroupService {

    @Autowired
    private GroupMapper groupMapper;

    @Autowired
    private GroupMemberMapper memberMapper;

    /**
     * 发群消息(写扩散)
     */
    public void sendGroupMessage(Long groupId, Message msg) {
        // 1. 存消息
        msg.setConversationType(ConversationType.GROUP);
        msg.setToId(groupId);
        storeService.save(msg);

        // 2. 推送给所有成员(除了发送者)
        List<Long> members = memberMapper.findMembersByGroupId(groupId);
        for (Long member : members) {
            if (member.equals(msg.getFromUserId())) continue;
            channelManager.sendToUser(member, msg);
        }

        // 3. Kafka 异步处理(已读回执、计数等)
        kafkaTemplate.send("group-message", msg);
    }
}

五、离线消息

5.1 离线消息存储

java
@Service
public class OfflineMessageService {

    @Autowired
    private RedisTemplate<String, String> redisTemplate;

    private static final int MAX_OFFLINE = 100;

    /**
     * 存离线消息
     */
    public void save(Long userId, Message msg) {
        String key = "offline:" + userId;
        redisTemplate.opsForList().leftPush(key, JSON.toJSONString(msg));
        redisTemplate.opsForList().trim(key, 0, MAX_OFFLINE - 1);

        // 设置过期 7 天
        redisTemplate.expire(key, 7, TimeUnit.DAYS);
    }

    /**
     * 拉取离线消息(用户上线时)
     */
    public List<Message> pull(Long userId) {
        String key = "offline:" + userId;
        List<String> list = redisTemplate.opsForList().range(key, 0, -1);
        if (list == null) return Collections.emptyList();

        // 拉取后清空
        redisTemplate.delete(key);

        return list.stream()
            .map(s -> JSON.parseObject(s, Message.class))
            .collect(Collectors.toList());
    }
}

5.2 多端同步

java
@Component
public class MultiDeviceSync {

    /** 用户的设备列表(通过 Redis Hash) */
    private static final String DEVICE_KEY = "devices:";

    @Autowired
    private RedisTemplate<String, String> redisTemplate;

    public void registerDevice(Long userId, String deviceId, String platform) {
        Map<String, String> info = new HashMap<>();
        info.put("deviceId", deviceId);
        info.put("platform", platform);
        info.put("loginAt", String.valueOf(System.currentTimeMillis()));
        redisTemplate.opsForHash().putAll(DEVICE_KEY + userId, info);
    }

    public void multiDeviceSync(Long userId, Message msg) {
        // 发送给所有在线设备
        Set<String> devices = redisTemplate.opsForZSet()
            .range(DEVICE_KEY + userId, 0, -1);

        if (devices != null) {
            for (String device : devices) {
                String channelId = device + ":" + userId;
                channelManager.sendToUserByChannel(channelId, msg);
            }
        }
    }
}

六、消息搜索

java
@Service
public class MessageSearchService {

    @Autowired
    private ElasticsearchClient esClient;

    @Autowired
    private ShardingJdbcTemplate jdbcTemplate;

    /**
     * 搜索消息(关键字)
     */
    public List<MessageVO> search(Long userId, String keyword, int page, int size) {
        // 1. 从 ES 搜索
        SearchRequest request = new SearchRequest("messages");
        SearchSourceBuilder builder = new SearchSourceBuilder()
            .query(QueryBuilders.boolQuery()
                .must(QueryTypes.termQuery("userId", userId))
                .must(QueryTypes.matchQuery("content", keyword)))
            .from(page * size)
            .size(size)
            .sort("timestamp", SortOrder.DESC);

        request.source(builder);

        SearchResponse response = esClient.search(request, RequestOptions.DEFAULT);

        List<String> msgIds = Arrays.stream(response.getHits().getHits())
            .map(hit -> hit.getId())
            .collect(Collectors.toList());

        // 2. 从 DB 加载详情
        if (msgIds.isEmpty()) return Collections.emptyList();
        List<MessageEntity> entities = jdbcTemplate.query(
            "SELECT * FROM t_message WHERE id IN (?)", msgIds
        );

        return entities.stream().map(this::toVO).collect(Collectors.toList());
    }

    /**
     * 数据同步(从 DB 到 ES)
     */
    @KafkaListener(topics = "message-index", groupId = "es-sync")
    public void indexMessage(Message msg) {
        IndexRequest request = new IndexRequest("messages");
        request.id(msg.getServerMsgId());
        request.source(JSON.toJSONString(msg), XContentType.JSON);
        esClient.index(request, RequestOptions.DEFAULT);
    }
}

七、推送集成

java
@Service
public class PushService {

    @Autowired
    private FcmClient fcmClient;

    @Autowired
    private JpushClient jpushClient;

    @Autowired
    private RedisTemplate<String, String> redisTemplate;

    /**
     * 离线推送
     */
    @KafkaListener(topics = "offline-push", groupId = "push")
    public void push(OfflinePushMessage msg) {
        // 1. 查用户设备列表
        List<DeviceInfo> devices = deviceService.getActiveDevices(msg.getUserId());
        if (devices.isEmpty()) return;

        for (DeviceInfo device : devices) {
            // 2. 选平台推送
            if ("iOS".equals(device.getPlatform())) {
                jpushClient.pushIos(device.getPushToken(), msg);
            } else if ("Android".equals(device.getPlatform())) {
                if ("domestic".equals(device.getRegion())) {
                    jpushClient.pushAndroid(device.getPushToken(), msg);
                } else {
                    fcmClient.push(device.getPushToken(), msg);
                }
            }
        }

        // 3. 写推送日志
        pushLogMapper.insert(new PushLog(msg, devices.size()));
    }
}

八、性能与监控

8.1 关键指标

yaml
连接层:
  - 在线人数
  - 消息吞吐(TPS)
  - 长连接建立/断开速率
  - 心跳延迟

业务层:
  - 消息发送成功率(> 99.9%)
  - 消息延迟 P99(< 100ms)
  - 离线消息送达率
  - 推送成功率

存储层:
  - 消息写入 QPS
  - DB 慢查询
  - Redis 命中率

8.2 消息堆积监控

java
@Component
public class MessageLagMonitor {

    @Autowired
    private KafkaAdmin kafkaAdmin;

    @Scheduled(fixedRate = 30000)
    public void check() {
        Map<String, TopicDescription> topics = kafkaAdmin.describeTopics(
            List.of("message-push", "message-sync")
        );

        for (TopicDescription topic : topics.values()) {
            for (TopicPartitionInfo partition : topic.partitions()) {
                long lag = getLag(partition);
                if (lag > 100000) {
                    alertService.send("消息堆积告警", partition.toString());
                }
            }
        }
    }
}

九、本章小结

模块关键
长连接Netty + Channel 管理
投递ACK + 重试 + 状态机
存储分库分表 + 冷热分离
离线Redis 队列 + 多端同步
推送FCM/极光 + 路由

动手练习

  1. 用 Netty 实现一个简单的 echo 服务
  2. 设计消息可靠投递的状态机
  3. 模拟离线消息的存取流程
  4. 调研你常用的 IM 客户端用了什么长连接方案

下一章:下一章:第 248 章:综合项目实战 - 网约车调度

本站基于 VitePress 构建 · 由 Codebook 团队维护