本文目录导读:

- 方案一:本地消息表 + 定时任务(最通用,适合中小项目)
- 方案二:基于RocketMQ的事务消息(需要安装RocketMQ扩展)
- 方案三:基于Redis队列 + 定时补偿(轻量级方案)
- 项目落地的关键考虑因素
- 推荐路线图
实现消息最终一致性在PHP项目中通常采用本地消息表+消息队列重试或事务性消息(如RocketMQ) 方案,考虑到PHP无常驻内存(传统FPM模式)的特性,以下提供两种最实用的落地方式。
本地消息表 + 定时任务(最通用,适合中小项目)
核心原理
- 业务操作与消息记录在同一个数据库事务中
- 定时任务扫描未确认的消息进行投递
- 消费端幂等处理
数据库表结构
-- 本地消息表 CREATE TABLE `message_queue` ( `id` bigint(20) NOT NULL AUTO_INCREMENT, `business_type` varchar(50) NOT NULL COMMENT '业务类型', `business_id` varchar(64) NOT NULL COMMENT '业务ID', `message_body` text NOT NULL COMMENT '消息体(JSON)', `status` tinyint(4) NOT NULL DEFAULT '0' COMMENT '0-待发送 1-已发送 2-已完成', `retry_count` int(11) NOT NULL DEFAULT '0' COMMENT '重试次数', `next_retry_time` datetime DEFAULT NULL COMMENT '下次重试时间', `create_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `update_time` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (`id`), KEY `idx_status_next_retry` (`status`, `next_retry_time`), KEY `idx_business` (`business_type`, `business_id`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;
PHP实现代码
业务写入(包含消息记录)
<?php
// 创建订单 + 写入消息(同一个事务)
class OrderService
{
private $db;
public function createOrder($orderData)
{
$this->db->beginTransaction();
try {
// 1. 业务操作 - 创建订单
$orderId = $this->insertOrder($orderData);
// 2. 同时写入本地消息表
$message = [
'order_id' => $orderId,
'amount' => $orderData['amount'],
'user_id' => $orderData['user_id']
];
$this->db->insert('message_queue', [
'business_type' => 'order_created',
'business_id' => $orderId,
'message_body' => json_encode($message, JSON_UNESCAPED_UNICODE),
'status' => 0, // 待发送
'next_retry_time' => date('Y-m-d H:i:s') // 立即执行
]);
$this->db->commit();
return $orderId;
} catch (\Exception $e) {
$this->db->rollback();
throw $e;
}
}
private function insertOrder($data)
{
// ... 订单插入逻辑
}
}
消息发送消费者(定时任务执行)
<?php
// 建议用crontab每分钟执行一次:* * * * * php /path/to/process_messages.php
class MessageProcessor
{
private $mq; // RabbitMQ或其他消息队列客户端
private $db;
private $maxRetryCount = 5;
public function processPendingMessages()
{
// 1. 获取待发送的消息(加锁防止并发)
$messages = $this->db->select(
"SELECT * FROM message_queue
WHERE status = 0
AND (next_retry_time IS NULL OR next_retry_time <= NOW())
AND retry_count < ?
ORDER BY id ASC
LIMIT 100
FOR UPDATE SKIP LOCKED",
[$this->maxRetryCount]
);
foreach ($messages as $msg) {
try {
// 2. 投递到真正的消息队列
$this->mq->publish(
$msg['business_type'],
$msg['message_body']
);
// 3. 更新状态为已发送
$this->db->update(
'message_queue',
['status' => 1],
['id' => $msg['id']]
);
} catch (\Exception $e) {
// 4. 更新重试次数和下次重试时间(指数退避)
$nextRetry = $this->calculateNextRetryTime($msg['retry_count'] + 1);
$this->db->update(
'message_queue',
[
'retry_count' => $msg['retry_count'] + 1,
'next_retry_time' => $nextRetry
],
['id' => $msg['id']]
);
// 记录错误日志
$this->logError($msg['id'], $e->getMessage());
}
}
}
private function calculateNextRetryTime($retryCount)
{
// 指数退避:1min, 2min, 4min, 8min, 16min...
$delay = pow(2, $retryCount - 1) * 60;
return date('Y-m-d H:i:s', time() + $delay);
}
}
消费端幂等处理
<?php
// 消息消费者
class OrderEventHandler
{
private $db;
public function handleOrderCreated($message)
{
$data = json_decode($message['body'], true);
// 幂等性校验:确保同一订单不会重复处理
$exists = $this->db->get(
"SELECT id FROM processed_messages
WHERE business_type = 'order_created'
AND business_id = ?",
[$data['order_id']]
);
if ($exists) {
// 已处理,直接ACK
return true;
}
$this->db->beginTransaction();
try {
// 执行实际业务逻辑(积分增加、库存扣减等)
// ...
// 记录处理成功的消息ID(幂等表)
$this->db->insert('processed_messages', [
'business_type' => 'order_created',
'business_id' => $data['order_id'],
'message_id' => $message['message_id'] ?? '',
'create_time' => date('Y-m-d H:i:s')
]);
$this->db->commit();
return true;
} catch (\Exception $e) {
$this->db->rollback();
throw $e; // 回滚,消息队列会重试
}
}
}
基于RocketMQ的事务消息(需要安装RocketMQ扩展)
RocketMQ原生支持事务消息,PHP可通过rocketmq-client-php扩展实现。
实现步骤
事务消息发送监听器
<?php
use RocketMQ\Producer;
use RocketMQ\TransactionListener;
use RocketMQ\Message;
class OrderTransactionListener implements TransactionListener
{
private $db;
// 执行本地事务
public function executeLocalTransaction(Message $msg)
{
$data = json_decode($msg->getBody(), true);
$this->db->beginTransaction();
try {
// 业务操作
$orderId = $this->createOrder($data);
// 事务成功
$this->db->commit();
return TransactionStatus::COMMIT_MESSAGE;
} catch (\Exception $e) {
$this->db->rollback();
return TransactionStatus::ROLLBACK_MESSAGE;
}
}
// 回查本地事务状态(RocketMQ补偿机制)
public function checkLocalTransaction(Message $msg)
{
$data = json_decode($msg->getBody(), true);
// 检查订单是否存在
$order = $this->db->get("SELECT id FROM orders WHERE id = ?", [$data['order_id']]);
return $order
? TransactionStatus::COMMIT_MESSAGE
: TransactionStatus::UNKNOWN;
}
}
// 发送事务消息
$producer = new Producer('127.0.0.1:9876');
$producer->start();
$msg = new Message('order_topic', json_encode(['order_id' => 123]));
$msg->setKeys('order_123');
$result = $producer->sendMessageInTransaction(
$msg, new OrderTransactionListener()
);
基于Redis队列 + 定时补偿(轻量级方案)
适合对一致性要求不那么极端的场景:
<?php
// 1. 业务操作时写入Redis List + DB记录
class RedisMessageService
{
private $redis;
private $db;
public function createOrderWithMessage($data)
{
$this->db->beginTransaction();
try {
$orderId = $this->insertOrder($data);
// 将消息写入Redis
$message = json_encode([
'type' => 'order_created',
'order_id' => $orderId,
'time' => date('Y-m-d H:i:s')
]);
$this->redis->lPush('order_queue', $message);
$this->db->commit();
return $orderId;
} catch (\Exception $e) {
$this->db->rollback();
throw $e;
}
}
}
// 2. 定时任务补偿(处理Redis可能丢失的情况)
// crontab每5分钟执行
function compensateMessages() {
// 扫描最近10分钟内未在Redis中处理的订单
$orders = $db->select(
"SELECT * FROM orders
WHERE status = 'pending'
AND create_time > DATE_SUB(NOW(), INTERVAL 10 MINUTE)
AND NOT EXISTS (
SELECT 1 FROM redis_processed
WHERE business_id = orders.id
)"
);
foreach ($orders as $order) {
// 重新投递到Redis
$redis->lPush('order_queue', json_encode($order));
}
}
项目落地的关键考虑因素
幂等性设计(最重要)
// 在消息表或单独幂等表中记录业务ID
// 消费前先检查是否已处理
function isProcessed($businessType, $businessId) {
return $this->db->get(
"SELECT 1 FROM idempotent_records
WHERE business_type = ? AND business_id = ?",
[$businessType, $businessId]
);
}
删除已处理消息的策略
- 方案A:保留7天后归档(推荐)
- 方案B:标记后保留30天用于排查问题
-- 归档旧数据 DELETE FROM message_queue WHERE status = 2 AND update_time < DATE_SUB(NOW(), INTERVAL 7 DAY);
监控与告警
// 在定时任务中增加监控
function monitorDeadLetters() {
$count = $this->db->get(
"SELECT COUNT(*) FROM message_queue
WHERE retry_count >= 5 AND status = 0"
);
if ($count > 10) {
alert("死信消息超过阈值:{$count}");
}
}
性能优化
- 本地消息表增加索引
- 定时任务分批处理(每次100条)
- 使用
SKIP LOCKED防止并发死锁(MySQL 8.0+) - 考虑用
Redis List+DB双写减少DB压力
推荐路线图
| 场景 | 推荐方案 |
|---|---|
| 小型项目、低并发 | 本地消息表 + crontab |
| 中等规模、需要高可靠 | 方案一升级版 + RabbitMQ/Kafka |
| 大型项目、强一致性 | RocketMQ事务消息 |
| 轻量级、允许少量丢失 | Redis队列 + 补偿 |
最终建议: 对于大多数PHP项目,方案一(本地消息表 + 定时任务 + RabbitMQ) 是最成熟、可控性最高的选择,它不需要依赖特定中间件,且完全符合最终一致性的理论要求,关键是做好幂等设计和失败补偿机制。