本文作者饭后咖啡,即时通讯网有修订和改动。
cover_small_opti.jpg (50.66 KB, 下载次数: 12)
下载附件 保存到相册
昨天 12:00 上传
1.png (23.97 KB, 下载次数: 10)
2.png (34.19 KB, 下载次数: 12)
昨天 12:01 上传
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亿
@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<OfflineMessage> 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<OfflineMessage> 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<OfflineMessage> 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<OfflineMessage> 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<OfflineMessage> 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);
@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<OfflineMessage> pullMessages(Long uid, int limit) { // 查询最近3天的分区 LocalDate today = LocalDate.now(); List<OfflineMessage> 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自动过期 } }
@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); } } }
3.png (19.57 KB, 下载次数: 12)
4.png (22.27 KB, 下载次数: 12)
@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<OfflineMessage> 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<OfflineMessage> all = new ArrayList<>(); all.addAll(hotMessages); all.addAll(coldMessages); return all; } }
5.png (24.44 KB, 下载次数: 11)
Redis + 定时落盘MySQL └── 实时消息存Redis(7天过期) └── 每天凌晨将Redis消息批量写入MySQL做永久备份
Redis热存 + HBase全量 + OSS大文件 └── Redis:存储最近3天消息(毫秒级拉取) └── HBase:存储所有消息(支持历史回溯) └── OSS:存储图片/视频等大文件 └── 统一API层屏蔽底层存储差异
自定义存储引擎 + 消息队列 + 分层存储 └── 基于RocksDB的本地存储(热数据) └── 分布式KV存储(温数据) └── 对象存储(冷数据/归档)
来源:即时通讯网 - 即时通讯开发者社区!
轻量级开源移动端即时通讯框架。
快速入门 / 性能 / 指南 / 提问
轻量级Web端即时通讯框架。
详细介绍 / 精编源码 / 手册教程
移动端实时音视频框架。
详细介绍 / 性能测试 / 安装体验
基于MobileIMSDK的移动IM系统。
详细介绍 / 产品截图 / 安装体验
一套产品级Web端IM系统。
详细介绍 / 产品截图 / 演示视频
一套纯血鸿蒙NEXT产品级IM系统。
详细介绍 / 产品截图 / 安装
精华主题数超过100个。
积极发起、参与各类话题的讨论等,主题、发帖内容较有价值。
连续任职达1年以上的合格正式版主
为论区做出突出贡献的开发者、版主等。
Copyright © 2014-2026 即时通讯网 - 即时通讯开发者社区 / 版本 V4.4
苏州网际时代信息科技有限公司 (苏ICP备16005070号-1)
Processed in 0.156250 second(s), 41 queries , Gzip On.