本文目录导读:

在PHP项目中实现客服聊天系统的后端消息存储,通常需要考虑数据结构、存储方案以及读写性能,以下是几种常见的实现思路和最佳实践:
数据库表设计
核心表结构
-- 会话表
CREATE TABLE `chat_session` (
`id` BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
`session_id` VARCHAR(64) NOT NULL COMMENT '会话唯一标识',
`customer_id` INT UNSIGNED NOT NULL COMMENT '客户ID',
`service_id` INT UNSIGNED DEFAULT NULL COMMENT '客服ID',
`status` TINYINT DEFAULT 1 COMMENT '状态: 1-进行中 2-已结束',
`created_at` DATETIME DEFAULT CURRENT_TIMESTAMP,
`updated_at` DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
UNIQUE KEY `uk_session_id` (`session_id`),
INDEX `idx_customer_id` (`customer_id`),
INDEX `idx_service_id` (`service_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
-- 消息表
CREATE TABLE `chat_message` (
`id` BIGINT UNSIGNED AUTO_INCREMENT PRIMARY KEY,
`session_id` VARCHAR(64) NOT NULL COMMENT '所属会话ID',
`sender_type` TINYINT NOT NULL COMMENT '发送者类型: 1-客户 2-客服 3-系统',
`sender_id` INT UNSIGNED NOT NULL COMMENT '发送者ID',
`message_type` TINYINT DEFAULT 1 COMMENT '消息类型: 1-文本 2-图片 3-文件 4-语音',
`content` TEXT NOT NULL COMMENT '消息内容',
`extra` JSON DEFAULT NULL COMMENT '扩展信息(图片URL、文件路径等)',
`read_status` TINYINT DEFAULT 0 COMMENT '已读状态: 0-未读 1-已读',
`created_at` DATETIME DEFAULT CURRENT_TIMESTAMP,
INDEX `idx_session_id` (`session_id`, `created_at`),
INDEX `idx_sender` (`sender_type`, `sender_id`),
INDEX `idx_read_status` (`session_id`, `read_status`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
PHP消息存储实现
基础消息存储类
<?php
class ChatMessageStorage
{
private $db;
public function __construct(PDO $db)
{
$this->db = $db;
}
/**
* 保存消息
*/
public function saveMessage(array $messageData): int
{
$sql = "INSERT INTO chat_message
(session_id, sender_type, sender_id, message_type, content, extra)
VALUES
(:session_id, :sender_type, :sender_id, :message_type, :content, :extra)";
$stmt = $this->db->prepare($sql);
$stmt->execute([
':session_id' => $messageData['session_id'],
':sender_type' => $messageData['sender_type'],
':sender_id' => $messageData['sender_id'],
':message_type' => $messageData['message_type'] ?? 1,
':content' => $messageData['content'],
':extra' => json_encode($messageData['extra'] ?? [])
]);
return (int)$this->db->lastInsertId();
}
/**
* 批量保存消息(用于缓存批量写入)
*/
public function batchSaveMessages(array $messages): array
{
$ids = [];
$this->db->beginTransaction();
try {
foreach ($messages as $message) {
$ids[] = $this->saveMessage($message);
}
$this->db->commit();
} catch (\Exception $e) {
$this->db->rollBack();
throw $e;
}
return $ids;
}
/**
* 获取会话消息历史
*/
public function getSessionMessages(string $sessionId, int $offset = 0, int $limit = 50): array
{
$sql = "SELECT * FROM chat_message
WHERE session_id = :session_id
ORDER BY created_at DESC
LIMIT :offset, :limit";
$stmt = $this->db->prepare($sql);
$stmt->bindValue(':session_id', $sessionId, PDO::PARAM_STR);
$stmt->bindValue(':offset', $offset, PDO::PARAM_INT);
$stmt->bindValue(':limit', $limit, PDO::PARAM_INT);
$stmt->execute();
return $stmt->fetchAll(PDO::FETCH_ASSOC);
}
/**
* 标记消息已读
*/
public function markAsRead(string $sessionId, int $senderId): bool
{
$sql = "UPDATE chat_message
SET read_status = 1
WHERE session_id = :session_id
AND sender_id = :sender_id
AND sender_type = 2"; // 标记客服发送的消息为已读
$stmt = $this->db->prepare($sql);
return $stmt->execute([
':session_id' => $sessionId,
':sender_id' => $senderId
]);
}
}
缓存层优化
对于高并发场景,建议引入Redis作为缓存层:
<?php
class ChatMessageCache
{
private $redis;
private $db;
private $messageQueue = [];
public function __construct(\Redis $redis, ChatMessageStorage $db)
{
$this->redis = $redis;
$this->db = $db;
}
/**
* 写入消息(先入缓存,异步持久化)
*/
public function writeMessage(array $message): int
{
// 1. 缓存最新消息到Redis
$cacheKey = "chat:session:{$message['session_id']}:messages";
$this->redis->lPush($cacheKey, json_encode($message));
$this->redis->lTrim($cacheKey, 0, 199); // 保留最近200条
// 2. 设置过期时间
$this->redis->expire($cacheKey, 3600);
// 3. 添加到批量写入队列
$message['id'] = 0; // 占位
$this->messageQueue[] = $message;
// 4. 达到阈值时批量写入数据库
if (count($this->messageQueue) >= 10) {
$this->flushToDatabase();
}
return $message['id'];
}
/**
* 批量写入数据库
*/
public function flushToDatabase(): void
{
if (empty($this->messageQueue)) {
return;
}
try {
$this->db->batchSaveMessages($this->messageQueue);
$this->messageQueue = [];
} catch (\Exception $e) {
// 记录日志,失败的消息可以重试
error_log("Chat message batch insert failed: " . $e->getMessage());
}
}
/**
* 获取消息历史(优先从缓存读取)
*/
public function getMessages(string $sessionId, int $offset = 0, int $limit = 50): array
{
$cacheKey = "chat:session:{$sessionId}:messages";
// 优先从缓存获取
if ($this->redis->exists($cacheKey)) {
$cached = $this->redis->lRange($cacheKey, $offset, $offset + $limit - 1);
$messages = array_map('json_decode', $cached);
if (count($messages) >= $limit) {
return $messages;
}
}
// 缓存不足时从数据库获取
$dbMessages = $this->db->getSessionMessages($sessionId, $offset, $limit);
// 回填缓存
foreach (array_reverse($dbMessages) as $msg) {
$this->redis->lPush($cacheKey, json_encode($msg));
}
return $dbMessages;
}
public function __destruct()
{
// 确保退出时写入剩余消息
$this->flushToDatabase();
}
}
WebSocket + 消息推送示例
<?php
// 使用 Ratchet 或 Swoole 实现WebSocket
use Ratchet\MessageComponentInterface;
use Ratchet\ConnectionInterface;
class ChatWebSocket implements MessageComponentInterface
{
private $clients;
private $messageStorage;
public function __construct(ChatMessageStorage $storage)
{
$this->clients = new \SplObjectStorage;
$this->messageStorage = $storage;
}
public function onOpen(ConnectionInterface $conn)
{
$this->clients->attach($conn);
}
public function onMessage(ConnectionInterface $from, $msg)
{
$data = json_decode($msg, true);
// 存储消息
$messageId = $this->messageStorage->saveMessage([
'session_id' => $data['session_id'],
'sender_type' => $data['sender_type'],
'sender_id' => $data['sender_id'],
'content' => $data['content']
]);
// 推送给接收方
foreach ($this->clients as $client) {
if ($client !== $from) {
$client->send(json_encode([
'type' => 'new_message',
'message_id' => $messageId,
'data' => $data
]));
}
}
}
public function onClose(ConnectionInterface $conn)
{
$this->clients->detach($conn);
}
public function onError(ConnectionInterface $conn, \Exception $e)
{
$conn->close();
}
}
消息分表策略
当消息量很大时,建议按时间分表:
<?php
class ChatMessageShard
{
/**
* 按月分表
*/
public function getTableName(string $sessionId, string $date = null): string
{
$date = $date ?? date('Y_m');
return "chat_message_{$date}";
}
/**
* 创建分表
*/
public function createMonthTable(string $yearMonth): void
{
$tableName = "chat_message_{$yearMonth}";
$sql = "CREATE TABLE IF NOT EXISTS `{$tableName}` LIKE `chat_message_template`";
// 执行建表SQL
}
}
性能优化建议
- 读写分离:主库写消息,从库读历史
- 消息队列:使用RabbitMQ或Redis List做异步写入
- 定期归档:将30天前的消息移到归档表
- 索引优化:根据查询模式创建合适的联合索引
- 连接池:使用长连接减少数据库连接开销
- 小型项目:直接使用MySQL + Redis缓存即可
- 中型项目:引入消息队列做异步写入,使用分表策略
- 大型项目:考虑使用专门的时序数据库(如ClickHouse)存储消息历史,实时消息用Redis Streams
选择合适的方案需要根据你的业务规模和并发量来决定,建议从简单方案开始,逐步演进。