如何用PHP项目实现消息最终一致性?

wen java案例 1

本文目录导读:

如何用PHP项目实现消息最终一致性?

  1. 方案一:本地消息表 + 定时任务(最通用,适合中小项目)
  2. 方案二:基于RocketMQ的事务消息(需要安装RocketMQ扩展)
  3. 方案三:基于Redis队列 + 定时补偿(轻量级方案)
  4. 项目落地的关键考虑因素
  5. 推荐路线图

实现消息最终一致性在PHP项目中通常采用本地消息表+消息队列重试事务性消息(如RocketMQ) 方案,考虑到PHP无常驻内存(传统FPM模式)的特性,以下提供两种最实用的落地方式。

本地消息表 + 定时任务(最通用,适合中小项目)

核心原理

  1. 业务操作与消息记录在同一个数据库事务中
  2. 定时任务扫描未确认的消息进行投递
  3. 消费端幂等处理

数据库表结构

-- 本地消息表
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) 是最成熟、可控性最高的选择,它不需要依赖特定中间件,且完全符合最终一致性的理论要求,关键是做好幂等设计和失败补偿机制。

抱歉,评论功能暂时关闭!