第 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/极光 + 路由 |
动手练习
- 用 Netty 实现一个简单的 echo 服务
- 设计消息可靠投递的状态机
- 模拟离线消息的存取流程
- 调研你常用的 IM 客户端用了什么长连接方案
下一章:下一章:第 248 章:综合项目实战 - 网约车调度