Jack Jiang

我的最新工程MobileIMSDK:http://git.oschina.net/jackjiang/MobileIMSDK
posts - 544, comments - 13, trackbacks - 0, articles - 1

本文作者饭后咖啡,有修订和改动。

1、引言

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

cover_small_opti

2、离线消息的业务特征和主流存储方案概览

离线消息的业务特征:

1

方案对比总览:

2

3、存储方案1:MySQL分库分表

表结构设计:

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;

    }

}

优点:

  • 1)简单可靠,事务支持;
  • 2)容易实现分页和删除。

缺点:

  • 1)写入吞吐有限(5k TPS);
  • 2)存储成本高,数据过期需手动清理。

4、存储方案2:Redis

数据结构设计:

@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);

    }

}

优点:

  • 1)读写性能极高(50k+ TPS);
  • 2)自动过期,无需清理任务;
  • 3)内存操作,延迟低(<5ms)。

缺点:

  • 1)存储成本高(内存是磁盘的10倍);
  • 2)消息量大时内存压力大。

5、存储方案3:HBase

表结构设计:

@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

    }

}

优点:

  • 1)海量存储能力(PB级);
  • 2)写入吞吐高(100k+ TPS);
  • 3)自动TTL过期。

缺点:

  • 1)查询灵活性差(只能按RowKey);
  • 2)运维复杂。

6、存储方案4:Cassandra

表结构设计:

-- 创建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自动过期

    }

}

优点:

  • 1)全球多活,跨地域部署;
  • 2)写入吞吐极高(200k+ TPS);
  • 3)无单点故障。

缺点:

  • 1)最终一致性,可能读到旧数据;
  • 2)查询能力受限。

7、存储方案5:RocketMQ

利用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);

        }

    }

}

优点:

  • 1)无需额外存储组件,复用MQ;
  • 2)天然支持消息顺序和重试;
  • 3)简单易实现。

缺点:

  • 1)长时间离线会堆积大量消息;
  • 2)消费能力受MQ性能限制。

8、存储方案6:对象存储 + Metadata

适用于超大附件和历史归档:

@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);

        }

    }

}

优点:

  • 1)存储成本极低;
  • 2)无限扩展能力;
  • 3)适合大文件。

缺点:

  • 1)访问延迟高(100ms+);
  • 2)不适合高频读写。

9、方案选型建议

按业务规模:

3

按业务特点:

4

10、混合架构实战

实际生产环境常采用混合架构:

@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;

    }

}

11、性能对比总结

5

12、最终建议

中小规模(用户<1000万):

Redis + 定时落盘MySQL

└── 实时消息存Redis(7天过期)

└── 每天凌晨将Redis消息批量写入MySQL做永久备份

大规模(用户>1000万):

Redis热存 + HBase全量 + OSS大文件

└── Redis:存储最近3天消息(毫秒级拉取)

└── HBase:存储所有消息(支持历史回溯)

└── OSS:存储图片/视频等大文件

└── 统一API层屏蔽底层存储差异

极端规模(十亿级用户):

自定义存储引擎 + 消息队列 + 分层存储

└── 基于RocksDB的本地存储(热数据)

└── 分布式KV存储(温数据) 

└── 对象存储(冷数据/归档)

选择方案时,建议从当前业务规模出发,预留2-3年的增长空间,不要过度设计。好的架构是演进出来的,不是一开始就完美的。

13、参考资料

[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聊天系统的消息可靠性(即消息不丢)

即时通讯技术学习:

- 移动端IM开发入门文章:《新手入门一篇就够:从零开发移动端IM

- 开源IM框架源码:https://github.com/JackJiang2011/MobileIMSDK备用地址点此

(本文同步发布于: http://www.52im.net/thread-4920-1-1.html



作者:Jack Jiang (点击作者姓名进入Github)
出处:http://www.52im.net/space-uid-1.html
交流:欢迎加入即时通讯开发交流群 215891622
讨论:http://www.52im.net/
Jack Jiang同时是【原创Java Swing外观工程BeautyEye】【轻量级移动端即时通讯框架MobileIMSDK】的作者,可前往下载交流。
本博文 欢迎转载,转载请注明出处(也可前往 我的52im.net 找到我)。


只有注册用户登录后才能发表评论。


网站导航:
 
Jack Jiang的 Mail: jb2011@163.com, 联系QQ: 413980957, 微信: hellojackjiang