
本文作者饭后咖啡,有修订和改动。
离线消息是IM系统的核心能力,当用户不在线时,消息必须可靠存储,待用户上线后完整推送。本文将对比主流的离线消息存储方案,帮你找到最适合业务场景的架构。

离线消息的业务特征:

方案对比总览:

表结构设计:
CREATE TABLE offline_message_$table ( id BIGINT NOT NULL AUTO_INCREMENT, receiver_uid BIGINT NOT NULL, -- 接收者用户ID sender_uid BIGINT NOT NULL, -- 发送者用户ID msg_id VARCHAR(64) NOT NULL, -- 消息全局ID msg_type TINYINT NOT NULL, -- 消息类型:文本/图片/语音 content TEXT NOT NULL, -- 消息内容 send_time DATETIME NOT NULL, -- 发送时间 status TINYINT DEFAULT 0, -- 状态:0未读,1已读,2撤回 PRIMARY KEY (id), UNIQUE KEY uk_msg_id (msg_id), KEY idx_receiver_time (receiver_uid, send_time) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4; -- 分片策略:按receiver_uid取模 -- 128张表,每表500万数据,总容量6.4亿
Java实现:
@Repository public class MySQLOfflineMessageStore { @Autowired private JdbcTemplate jdbcTemplate; private static final int TABLE_COUNT = 128; /** * 存储离线消息 */ public void saveMessage(OfflineMessage msg) { String table = getTableName(msg.getReceiverUid()); String sql = "INSERT INTO " + table + "(receiver_uid, sender_uid, msg_id, msg_type, content, send_time) " + "VALUES (?, ?, ?, ?, ?, ?)"; jdbcTemplate.update(sql, msg.getReceiverUid(), msg.getSenderUid(), msg.getMsgId(), msg.getMsgType(), msg.getContent(), msg.getSendTime() ); } /** * 拉取离线消息(分页) */ public List pullMessages(Long uid, Long lastMsgId, int limit) { String table = getTableName(uid); String sql = "SELECT * FROM " + table + " WHERE receiver_uid = ? AND id > ? " + " ORDER BY id ASC LIMIT ?"; return jdbcTemplate.query(sql, new Object[]{uid, lastMsgId, limit}, new BeanPropertyRowMapper<>(OfflineMessage.class) ); } /** * 批量删除已拉取的消息 */ public void deletePulledMessages(Long uid, List<Long> msgIds) { String table = getTableName(uid); String sql = "DELETE FROM " + table + " WHERE receiver_uid = ? AND id IN (" + msgIds.stream().map(String::valueOf).collect(Collectors.joining(",")) + ")"; jdbcTemplate.update(sql, uid); } private String getTableName(Long uid) { int index = (int) (uid % TABLE_COUNT); return "offline_message_" + index; } }
优点:
缺点:
数据结构设计:
@Component public class RedisOfflineMessageStore { @Autowired private RedisTemplate<String, Object> redisTemplate; private static final String OFFLINE_MSG_KEY = "offline:msg:"; private static final String OFFLINE_QUEUE_KEY = "offline:queue:"; private static final Duration MESSAGE_TTL = Duration.ofDays(7); /** * 存储离线消息 * 使用List结构,每个用户一个队列 */ public void saveMessage(OfflineMessage msg) { String msgKey = OFFLINE_MSG_KEY + msg.getMsgId(); String queueKey = OFFLINE_QUEUE_KEY + msg.getReceiverUid(); // 1. 存储消息内容(7天过期) redisTemplate.opsForValue().set(msgKey, msg, MESSAGE_TTL); // 2. 将消息ID推入用户队列 redisTemplate.opsForList().rightPush(queueKey, msg.getMsgId()); redisTemplate.expire(queueKey, MESSAGE_TTL); } /** * 拉取离线消息 */ public List pullMessages(Long uid, int limit) { String queueKey = OFFLINE_QUEUE_KEY + uid; // 从队列左侧批量弹出消息ID List<Object> msgIds = redisTemplate.opsForList() .leftPop(queueKey, limit); if (msgIds == null || msgIds.isEmpty()) { return Collections.emptyList(); } // 批量获取消息内容 List messages = new ArrayList<>(); for (Object msgId : msgIds) { String msgKey = OFFLINE_MSG_KEY + msgId; OfflineMessage msg = (OfflineMessage) redisTemplate.opsForValue().get(msgKey); if (msg != null) { messages.add(msg); } // 删除已拉取的消息内容 redisTemplate.delete(msgKey); } return messages; } /** * 使用SortedSet按时间排序(适用于需要保留历史的场景) */ public void saveMessageWithTime(OfflineMessage msg) { String msgKey = OFFLINE_MSG_KEY + msg.getMsgId(); String sortedSetKey = "offline:zset:" + msg.getReceiverUid(); // 存储消息 redisTemplate.opsForValue().set(msgKey, msg, MESSAGE_TTL); // 使用时间戳作为score redisTemplate.opsForZSet().add( sortedSetKey, msg.getMsgId(), msg.getSendTime().getTime() ); redisTemplate.expire(sortedSetKey, MESSAGE_TTL); } }
优点:
缺点:
表结构设计:
@Configuration public class HBaseOfflineMessageConfig { @Bean public HBaseTemplate hbaseTemplate() { return new HBaseTemplate(connection); } /** * 创建表:offline_message * RowKey设计:receiver_uid + (Long.MAX_VALUE - timestamp) * 保证同一用户的消息按时间倒序排列 */ public void createTable() { HBaseAdmin admin = (HBaseAdmin) connection.getAdmin(); TableName tableName = TableName.valueOf("offline_message"); if (!admin.tableExists(tableName)) { TableDescriptor desc = TableDescriptorBuilder.newBuilder(tableName) .setColumnFamily(ColumnFamilyDescriptorBuilder .newBuilder(Bytes.toBytes("info")) .setMaxVersions(1) .setTimeToLive(7 * 24 * 60 * 60) // 7天过期 .build()) .build(); admin.createTable(desc); } } } @Repository public class HBaseOfflineMessageStore { @Autowired private HBaseTemplate hbaseTemplate; private static final String TABLE_NAME = "offline_message"; private static final String FAMILY = "info"; /** * 存储消息 */ public void saveMessage(OfflineMessage msg) { // RowKey: receiver_uid + (Long.MAX_VALUE - timestamp) long reverseTime = Long.MAX_VALUE - msg.getSendTime().getTime(); byte[] rowKey = Bytes.add( Bytes.toBytes(msg.getReceiverUid()), Bytes.toBytes(reverseTime) ); hbaseTemplate.put(TABLE_NAME, rowKey, FAMILY, (put) -> { put.addColumn(Bytes.toBytes(FAMILY), Bytes.toBytes("sender"), Bytes.toBytes(msg.getSenderUid())); put.addColumn(Bytes.toBytes(FAMILY), Bytes.toBytes("msg_id"), Bytes.toBytes(msg.getMsgId())); put.addColumn(Bytes.toBytes(FAMILY), Bytes.toBytes("msg_type"), Bytes.toBytes(msg.getMsgType())); put.addColumn(Bytes.toBytes(FAMILY), Bytes.toBytes("content"), Bytes.toBytes(msg.getContent())); put.addColumn(Bytes.toBytes(FAMILY), Bytes.toBytes("send_time"), Bytes.toBytes(msg.getSendTime().getTime())); }); } /** * 拉取离线消息 */ public List pullMessages(Long uid, long lastTimestamp, int limit) { byte[] startRow; byte[] endRow; if (lastTimestamp > 0) { // 拉取比lastTimestamp更早的消息 startRow = Bytes.add( Bytes.toBytes(uid), Bytes.toBytes(Long.MAX_VALUE - lastTimestamp) ); endRow = Bytes.add(Bytes.toBytes(uid), Bytes.toBytes(0L)); } else { // 首次拉取,从最新开始 startRow = Bytes.add(Bytes.toBytes(uid), Bytes.toBytes(Long.MAX_VALUE)); endRow = Bytes.add(Bytes.toBytes(uid), Bytes.toBytes(0L)); } Scan scan = new Scan() .withStartRow(startRow) .withStopRow(endRow) .setMaxResultSize(limit) .setReversed(true); // 倒序 List messages = new ArrayList<>(); hbaseTemplate.find(TABLE_NAME, scan, (result) -> { OfflineMessage msg = new OfflineMessage(); msg.setReceiverUid(uid); msg.setSenderUid(Bytes.toLong( result.getValue(FAMILY, "sender"))); msg.setMsgId(Bytes.toString( result.getValue(FAMILY, "msg_id"))); msg.setMsgType(Bytes.toInt( result.getValue(FAMILY, "msg_type"))); msg.setContent(Bytes.toString( result.getValue(FAMILY, "content"))); msg.setSendTime(new Date(Bytes.toLong( result.getValue(FAMILY, "send_time")))); messages.add(msg); return false; }); return messages; } /** * 删除已拉取的消息(由HBase TTL自动处理,无需显式删除) */ public void deleteMessages(Long uid, List<String> msgIds) { // HBase通常依赖TTL自动过期 // 如需立即删除,可批量删除指定RowKey } }
优点:
缺点:
表结构设计:
-- 创建Keyspace CREATE KEYSPACE IF NOT EXISTS im WITH replication = {'class': 'NetworkTopologyStrategy', 'dc1': '3'}; -- 创建离线消息表 CREATE TABLE im.offline_message ( receiver_uid bigint, bucket_time timestamp, -- 时间桶,按天分区 send_time timestamp, msg_id uuid, sender_uid bigint, msg_type int, content text, PRIMARY KEY ((receiver_uid, bucket_time), send_time, msg_id) ) WITH CLUSTERING ORDER BY (send_time DESC) AND default_time_to_live = 604800; -- 7天过期 -- 创建索引(用于清理等操作) CREATE INDEX ON im.offline_message (sender_uid);
Java实现:
@Repository public class CassandraOfflineMessageStore { @Autowired private CassandraTemplate cassandraTemplate; private static final int BUCKET_SIZE_DAYS = 1; /** * 存储消息 */ public void saveMessage(OfflineMessage msg) { // 计算时间桶(按天分区) LocalDate bucketDate = msg.getSendTime() .toInstant() .atZone(ZoneId.systemDefault()) .toLocalDate(); OfflineMessageEntity entity = new OfflineMessageEntity(); entity.setReceiverUid(msg.getReceiverUid()); entity.setBucketTime(bucketDate); entity.setSendTime(msg.getSendTime()); entity.setMsgId(UUID.randomUUID()); entity.setSenderUid(msg.getSenderUid()); entity.setMsgType(msg.getMsgType()); entity.setContent(msg.getContent()); cassandraTemplate.insert(entity); } /** * 拉取消息 */ public List pullMessages(Long uid, int limit) { // 查询最近3天的分区 LocalDate today = LocalDate.now(); List messages = new ArrayList<>(); for (int i = 0; i < 3; i++) { LocalDate bucketDate = today.minusDays(i); Select select = QueryBuilder.select() .from("offline_message") .where(QueryBuilder.eq("receiver_uid", uid)) .and(QueryBuilder.eq("bucket_time", bucketDate)) .limit(limit - messages.size()); List<OfflineMessageEntity> entities = cassandraTemplate.select(select, OfflineMessageEntity.class); for (OfflineMessageEntity entity : entities) { messages.add(convert(entity)); } if (messages.size() >= limit) { break; } } return messages; } /** * 删除消息(可选) */ public void deleteMessages(Long uid, List<String> msgIds) { // 根据实际情况实现删除 // 通常依赖TTL自动过期 } }
优点:
缺点:
利用RocketMQ的消息重试机制实现离线消息:
@Component public class RocketMQOfflineMessageStore { @Autowired private RocketMQTemplate rocketMQTemplate; @Autowired private UserOnlineStatusService onlineStatusService; private static final String TOPIC_P2P = "im-p2p"; private static final int MAX_RETRY_TIMES = 7; // 最多重试7天 /** * 发送点对点消息 */ public void sendMessage(OfflineMessage msg) { Message<OfflineMessage> message = MessageBuilder .withPayload(msg) .setHeader("receiver_uid", msg.getReceiverUid()) .setHeader("retry_times", 0) .build(); // 发送到用户专属队列 rocketMQTemplate.syncSend( TOPIC_P2P + ":" + msg.getReceiverUid(), message, TIMEOUT ); } /** * 消费消息(消费者) */ @RocketMQMessageListener( topic = TOPIC_P2P, selectorExpression = "receiver_uid", consumerGroup = "im-consumer" ) public class IMConsumer implements RocketMQListener<OfflineMessage> { @Override public void onMessage(OfflineMessage message) { Long receiverUid = message.getReceiverUid(); // 检查用户是否在线 if (onlineStatusService.isOnline(receiverUid)) { // 在线,直接推送 pushToUser(receiverUid, message); } else { // 离线,抛出异常触发重试 throw new RuntimeException("User offline, retry later"); } } } /** * 配置重试策略 */ @Bean public DefaultMQPushConsumer imConsumer() { DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("im-group"); // 设置重试次数和延迟级别 // 延迟级别:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h consumer.setMaxReconsumeTimes(MAX_RETRY_TIMES); // 配置消息超时后进入死信队列 consumer.setConsumeTimeout(15); // 15分钟 return consumer; } /** * 处理死信消息(超过重试次数) */ @RocketMQMessageListener( topic = "%DLQ%im-consumer", consumerGroup = "im-dlq-processor" ) public class DeadLetterProcessor implements RocketMQListener<OfflineMessage> { @Override public void onMessage(OfflineMessage message) { // 消息重试7次后仍失败,转为持久化存储 saveToLongTermStorage(message); } } }
优点:
缺点:
适用于超大附件和历史归档:
@Component public class ObjectStorageOfflineMessageStore { @Autowired private OSSClient ossClient; @Autowired private JdbcTemplate metaJdbcTemplate; private static final String BUCKET_NAME = "im-messages"; private static final String OSS_ENDPOINT = "https://oss.aliyuncs.com"; /** * 存储大消息(图片、视频、文件等) */ public void saveLargeMessage(OfflineMessage msg, byte[] data) { // 1. 生成OSS路径 String objectKey = String.format("msg/%d/%d/%s.dat", msg.getReceiverUid() / 10000, msg.getReceiverUid(), msg.getMsgId()); // 2. 上传到OSS ossClient.putObject(BUCKET_NAME, objectKey, new ByteArrayInputStream(data)); // 3. 存储元数据到MySQL String sql = "INSERT INTO message_metadata " + "(msg_id, receiver_uid, sender_uid, msg_type, " + "object_key, file_size, send_time) " + "VALUES (?, ?, ?, ?, ?, ?, ?)"; metaJdbcTemplate.update(sql, msg.getMsgId(), msg.getReceiverUid(), msg.getSenderUid(), msg.getMsgType(), objectKey, data.length, msg.getSendTime() ); } /** * 获取消息内容 */ public byte[] getMessageContent(String msgId) { // 1. 查询元数据 String sql = "SELECT object_key FROM message_metadata " + "WHERE msg_id = ?"; String objectKey = metaJdbcTemplate.queryForObject( sql, String.class, msgId); // 2. 从OSS下载 OSSObject ossObject = ossClient.getObject( BUCKET_NAME, objectKey); try (InputStream in = ossObject.getObjectContent()) { return IOUtils.toByteArray(in); } catch (IOException e) { throw new RuntimeException(e); } } }
优点:
缺点:
按业务规模:

按业务特点:

实际生产环境常采用混合架构:
@Component public class HybridOfflineMessageStore { @Autowired private RedisOfflineMessageStore redisStore; @Autowired private HBaseOfflineMessageStore hbaseStore; private static final int HOT_DAYS = 3; // 热数据保留3天 /** * 存储消息:同时写入热存储和冷存储 */ public void saveMessage(OfflineMessage msg) { // 写入Redis(热数据) redisStore.saveMessage(msg); // 异步写入HBase(全量数据) CompletableFuture.runAsync(() -> { hbaseStore.saveMessage(msg); }); } /** * 拉取消息:优先从热存储读取 */ public List pullMessages(Long uid, Long lastMsgId, int limit) { // 1. 先从Redis拉取最近消息 List<OfflineMessage> hotMessages = redisStore.pullMessages(uid, limit); if (hotMessages.size() >= limit) { return hotMessages; } // 2. Redis不够,再从HBase拉取历史 int remaining = limit - hotMessages.size(); List<OfflineMessage> coldMessages = hbaseStore.pullMessages( uid, getLastTimestamp(hotMessages), remaining); // 3. 合并结果 List all = new ArrayList<>(); all.addAll(hotMessages); all.addAll(coldMessages); return all; } }

中小规模(用户<1000万):
Redis + 定时落盘MySQL └── 实时消息存Redis(7天过期) └── 每天凌晨将Redis消息批量写入MySQL做永久备份
大规模(用户>1000万):
Redis热存 + HBase全量 + OSS大文件 └── Redis:存储最近3天消息(毫秒级拉取) └── HBase:存储所有消息(支持历史回溯) └── OSS:存储图片/视频等大文件 └── 统一API层屏蔽底层存储差异
极端规模(十亿级用户):
自定义存储引擎 + 消息队列 + 分层存储 └── 基于RocksDB的本地存储(热数据) └── 分布式KV存储(温数据) └── 对象存储(冷数据/归档)
选择方案时,建议从当前业务规模出发,预留2-3年的增长空间,不要过度设计。好的架构是演进出来的,不是一开始就完美的。
[1] 零基础IM开发入门(一):什么是IM聊天系统?
[2] 一套海量在线用户的移动端IM架构设计实践分享(含详细图文)
[3] 某信团队分享:来看看微信十年前的IM消息收发架构,你做到了吗
[4] 如何保障分布式IM聊天系统的消息可靠性(即消息不丢)
[5] IM消息送达保证机制实现(二):保证离线消息的可靠投递
[6] IM群聊消息如此复杂,如何保证不丢不重?
[7] IM开发干货分享:如何优雅的实现大量离线消息的可靠投递
[8] 一套亿级用户的IM架构技术干货(下篇):可靠性、有序性、弱网优化等
[9] 阿里IM技术分享(六):亿级IM消息系统的离线推送到达率优化
[10] 阿里IM技术分享(七):IM的在线、离线聊天数据同步机制优化实践
[11] 如何保障分布式IM聊天系统的消息可靠性(即消息不丢)
(本文同步发布于: http://www.52im.net/thread-4920-1-1.html)
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。