本文目录导读:

我来详细说明PHP消息漫游的实现方案:
消息漫游概述
消息漫游是指用户在不同设备上登录时,可以同步查看历史聊天记录的功能。
核心架构设计
1 数据存储架构
// 数据库表设计
CREATE TABLE messages (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
msg_id VARCHAR(64) UNIQUE, -- 全局唯一消息ID
from_user_id INT NOT NULL,
to_user_id INT NOT NULL,
conversation_id VARCHAR(64), -- 会话ID
msg_type TINYINT, -- 消息类型
content MEDIUMTEXT, -- 消息内容
client_msg_id VARCHAR(64), -- 客户端消息ID
status TINYINT DEFAULT 1, -- 消息状态
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
INDEX idx_conversation_time (conversation_id, created_at),
INDEX idx_user_time (from_user_id, created_at)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 用户设备表
CREATE TABLE user_devices (
id INT PRIMARY KEY AUTO_INCREMENT,
user_id INT NOT NULL,
device_id VARCHAR(64),
device_type VARCHAR(20),
last_login_time TIMESTAMP,
UNIQUE KEY uk_user_device (user_id, device_id)
) ENGINE=InnoDB;
2 消息ID生成器
<?php
class MessageIdGenerator {
private $workerId;
public function __construct($workerId) {
$this->workerId = $workerId;
}
public function generate() {
$timestamp = microtime(true) * 1000;
// 使用Redis生成自增序列
$redis = Redis::getInstance();
$seq = $redis->incr('message:seq');
// 组装唯一ID
$msgId = sprintf('%s%04d%06d',
date('YmdHis'),
$this->workerId,
$seq % 1000000
);
return $msgId;
}
}
消息存储策略
1 热数据缓存
<?php
class MessageCache {
private $redis;
private $cachePrefix = 'msg:';
public function __construct() {
$this->redis = new Redis();
}
/**
* 缓存最近消息
*/
public function cacheRecentMessages($conversationId, $messages) {
$key = $this->cachePrefix . $conversationId;
foreach ($messages as $msg) {
$this->redis->rPush($key, json_encode($msg));
}
// 设置过期时间(7天)
$this->redis->expire($key, 7 * 86400);
// 只保留最近100条
$this->redis->lTrim($key, -100, -1);
}
/**
* 获取缓存消息
*/
public function getCachedMessages($conversationId, $start, $end) {
$key = $this->cachePrefix . $conversationId;
$messages = $this->redis->lRange($key, $start, $end);
return array_map(function($msg) {
return json_decode($msg, true);
}, $messages);
}
}
2 冷数据持久化
<?php
class MessageRepository {
private $db;
private $cache;
/**
* 同步消息到数据库
*/
public function syncToDatabase($message) {
$sql = "INSERT INTO messages (msg_id, from_user_id, to_user_id,
conversation_id, msg_type, content, client_msg_id)
VALUES (?, ?, ?, ?, ?, ?, ?)";
$stmt = $this->db->prepare($sql);
$stmt->execute([
$message['msg_id'],
$message['from_user_id'],
$message['to_user_id'],
$message['conversation_id'],
$message['msg_type'],
$message['content'],
$message['client_msg_id']
]);
}
/**
* 获取历史消息(分页)
*/
public function getHistoryMessages($userId, $conversationId, $page, $pageSize) {
$offset = ($page - 1) * $pageSize;
$sql = "SELECT * FROM messages
WHERE conversation_id = ?
AND id < ?
ORDER BY created_at DESC
LIMIT ? OFFSET ?";
$stmt = $this->db->prepare($sql);
$stmt->execute([$conversationId, $offset + $pageSize, $pageSize, $offset]);
return $stmt->fetchAll(PDO::FETCH_ASSOC);
}
}
消息同步机制
1 增量同步
<?php
class MessageSync {
private $userId;
private $syncService;
/**
* 获取增量消息
*/
public function getIncrementalMessages($userId, $lastSyncTime) {
$messages = [];
// 从数据库获取上次同步后的新消息
$sql = "SELECT * FROM messages
WHERE (from_user_id = ? OR to_user_id = ?)
AND created_at > ?
ORDER BY created_at ASC
LIMIT 100";
$stmt = $this->db->prepare($sql);
$stmt->execute([$userId, $userId, $lastSyncTime]);
$newMessages = $stmt->fetchAll();
// 更新同步游标
$this->saveSyncCursor($userId, time());
return $newMessages;
}
/**
* 保存同步游标
*/
private function saveSyncCursor($userId, $cursor) {
$key = "sync:cursor:{$userId}";
Redis::getInstance()->set($key, $cursor);
}
}
2 设备同步
<?php
class DeviceSync {
/**
* 推送消息到用户所有在线设备
*/
public function pushToAllDevices($userId, $message) {
// 获取用户所有设备
$devices = $this->getUserDevices($userId);
foreach ($devices as $device) {
if ($device['online']) {
$this->pushToDevice($device['device_id'], $message);
}
}
}
/**
* 心跳机制维护连接状态
*/
public function heartbeat($userId, $deviceId) {
$key = "device:online:{$userId}:{$deviceId}";
Redis::getInstance()->setex($key, 60, 'online');
}
}
消息漫游API
<?php
class MessageRoamingAPI {
/**
* 拉取历史消息
*/
public function getRoamingMessages($request) {
$userId = $this->getCurrentUserId();
$conversationId = $request['conversation_id'];
$page = isset($request['page']) ? (int)$request['page'] : 1;
$pageSize = isset($request['page_size']) ? min($request['page_size'], 100) : 20;
$lastMsgId = isset($request['last_msg_id']) ? $request['last_msg_id'] : '';
// 1. 先从缓存获取
$cacheMessages = $this->getFromCache($conversationId, $page, $pageSize);
// 2. 缓存不够,从数据库获取
if (count($cacheMessages) < $pageSize) {
$dbMessages = $this->getFromDatabase(
$userId,
$conversationId,
$page,
$pageSize
);
// 合并消息并去重
$messages = $this->mergeMessages($cacheMessages, $dbMessages);
} else {
$messages = $cacheMessages;
}
// 3. 标记已同步
$this->updateSyncStatus($userId, $conversationId, $lastMsgId);
return [
'code' => 0,
'data' => [
'messages' => $messages,
'has_more' => $this->hasMoreMessages($conversationId, $page, $pageSize)
]
];
}
/**
* 同步全部会话
*/
public function syncAllConversations($request) {
$userId = $this->getCurrentUserId();
$lastSyncTime = $request['last_sync_time'] ?? 0;
// 获取用户所有会话
$conversations = $this->getUserConversations($userId);
// 获取每个会话的最新消息
$result = [];
foreach ($conversations as $conv) {
$messages = $this->getIncrementalMessages(
$conv['conversation_id'],
$lastSyncTime
);
$result[] = [
'conversation_id' => $conv['conversation_id'],
'messages' => $messages
];
}
return $result;
}
}
性能优化
1 分库分表
<?php
class MessageSharding {
/**
* 根据会话ID计算分表
*/
public function getTableName($conversationId) {
$hash = crc32($conversationId);
$tableIndex = $hash % 100;
return "messages_{$tableIndex}";
}
/**
* 根据用户ID计算分库
*/
public function getDatabaseName($userId) {
$dbIndex = $userId % 10;
return "chat_db_{$dbIndex}";
}
}
2 消息压缩
<?php
class MessageCompression {
/**
* 压缩消息内容
*/
public function compressMessage($message) {
// 去掉无效字符
$cleaned = preg_replace('/\s+/', ' ', $message['content']);
// 压缩
return gzcompress($cleaned, 9);
}
/**
* 解压消息
*/
public function decompressMessage($compressed) {
return gzuncompress($compressed);
}
}
安全与可靠性
1 消息加密
<?php
class MessageEncryption {
private $key;
/**
* 加密消息
*/
public function encrypt($data) {
$method = 'AES-256-CBC';
$iv = openssl_random_pseudo_bytes(16);
$encrypted = openssl_encrypt($data, $method, $this->key, 0, $iv);
return base64_encode($iv . $encrypted);
}
/**
* 解密消息
*/
public function decrypt($data) {
$method = 'AES-256-CBC';
$raw = base64_decode($data);
$iv = substr($raw, 0, 16);
$encrypted = substr($raw, 16);
return openssl_decrypt($encrypted, $method, $this->key, 0, $iv);
}
}
2 消息去重
<?php
class MessageDeduplicator {
private $redis;
/**
* 检测重复消息
*/
public function isDuplicate($clientMsgId, $userId) {
$key = "msg:dedup:{$userId}:{$clientMsgId}";
// SETNX命令
$result = $this->redis->setnx($key, '1');
if ($result) {
// 设置过期时间,防止内存泄漏
$this->redis->expire($key, 24 * 3600);
return false; // 不是重复消息
}
return true; // 是重复消息
}
}
完整使用示例
<?php
// 发送消息
class MessageManager {
public function sendMessage($fromUserId, $toUserId, $content) {
// 1. 生成消息ID
$msgId = (new MessageIdGenerator(1))->generate();
// 2. 构建消息对象
$message = [
'msg_id' => $msgId,
'from_user_id' => $fromUserId,
'to_user_id' => $toUserId,
'conversation_id' => $this->getConversationId($fromUserId, $toUserId),
'msg_type' => 1, // 文本消息
'content' => $content,
'created_at' => time()
];
// 3. 同步到缓存和数据库
$this->cache->cacheRecentMessages($message['conversation_id'], [$message]);
$this->repository->syncToDatabase($message);
// 4. 推送通知
$this->notifyUser($toUserId, $message);
return $message;
}
/**
* 获取会话ID
*/
private function getConversationId($user1, $user2) {
$ids = [$user1, $user2];
sort($ids);
return implode('_', $ids);
}
}
// 初始化
$messageManager = new MessageManager();
// 发送消息
$result = $messageManager->sendMessage(1001, 1002, "你好,这是测试消息");
// 用户1002在不同设备拉取历史消息
$roamingApi = new MessageRoamingAPI();
$history = $roamingApi->getRoamingMessages([
'conversation_id' => '1001_1002',
'page' => 1,
'page_size' => 20
]);
监控与维护
<?php
class MessageMonitor {
/**
* 监控消息同步延迟
*/
public function checkSyncDelay() {
$redis = Redis::getInstance();
$lastSyncTime = $redis->get('last:sync:time') ?? 0;
$currentDelay = time() - $lastSyncTime;
if ($currentDelay > 300) {
// 发送告警
$this->sendAlert('消息同步延迟过高', $currentDelay);
}
return $currentDelay;
}
/**
* 监控存储容量
*/
public function checkStorageCapacity() {
$db = Database::getInstance();
$result = $db->query("SHOW TABLE STATUS LIKE 'messages_%'")->fetchAll();
$totalSize = 0;
foreach ($result as $table) {
$totalSize += $table['Data_length'] + $table['Index_length'];
}
return $totalSize / (1024 * 1024); // MB
}
}
最佳实践建议
- 消息有效期设定:根据业务需求设置消息存储期限(如180天)
- 分级存储:热数据放Redis,温数据放MySQL,冷数据放OSS
- 异步同步:消息同步采用异步方式,避免阻塞主流程
- 失败重试机制:启用可靠的消息投递
- 监控告警:实时监控存储容量、同步延迟等指标
- 扩展性设计:做好分库分表和集群扩展规划
这套方案能够满足绝大多数消息漫游需求,根据实际业务场景可以适当调整。