PHP 事务消息怎么实现

wen PHP项目 2

本文目录导读:

PHP 事务消息怎么实现

  1. 核心问题
  2. 常用实现方案
  3. 最佳实践建议

在 PHP 中实现事务消息(Transactional Messaging),主要涉及数据库事务消息队列的一致性保证,下面我会从核心原理到具体实现,逐层展开说明。

核心问题

事务消息要解决的核心问题是分布式事务:数据库操作和消息发送必须同时成功或同时失败。

  • 用户下单(写数据库)后必须发送订单消息
  • 如果数据库写入成功但消息发送失败,会导致数据不一致

常用实现方案

本地消息表(Local Message Table)

这是最经典、最可靠的方案。

<?php
// 1. 在主业务数据库中创建消息表
class OrderService
{
    private $pdo;
    private $mq;
    public function __construct(PDO $pdo, MessageQueue $mq)
    {
        $this->pdo = $pdo;
        $this->mq = $mq;
    }
    public function createOrder($orderData)
    {
        // 开启事务
        $this->pdo->beginTransaction();
        try {
            // 1. 插入订单数据
            $orderId = $this->insertOrder($orderData);
            // 2. 同时插入本地消息表
            $this->insertLocalMessage($orderId, [
                'event' => 'ORDER_CREATED',
                'data' => $orderData
            ]);
            // 提交事务 - 订单和消息同时落库
            $this->pdo->commit();
            // 3. 事务提交后,异步发送消息
            $this->sendMessageAfterCommit($orderId);
            return $orderId;
        } catch (Exception $e) {
            $this->pdo->rollBack();
            throw $e;
        }
    }
    private function insertLocalMessage($orderId, $messageData)
    {
        $sql = "INSERT INTO message_queue (order_id, message_data, status, create_time) 
                VALUES (?, ?, 'PENDING', NOW())";
        $stmt = $this->pdo->prepare($sql);
        $stmt->execute([$orderId, json_encode($messageData)]);
    }
    private function sendMessageAfterCommit($orderId)
    {
        // 通过中间件或事件驱动异步处理
        event(new OrderCreated($orderId));
    }
}

基于 PDO 的事务回调

利用 PDO 的提交后处理机制:

<?php
class TransactionalMessenger
{
    private $pdo;
    private $pendingMessages = [];
    public function __construct(PDO $pdo)
    {
        $this->pdo = $pdo;
    }
    public function sendAfterCommit($event, $data)
    {
        // 将消息暂存到内存
        $this->pendingMessages[] = compact('event', 'data');
    }
    public function executeInTransaction(callable $callback)
    {
        $this->pdo->beginTransaction();
        try {
            $result = $callback($this);
            $this->pdo->commit();
            // 事务提交成功后发送消息
            $this->flushMessages();
            return $result;
        } catch (Exception $e) {
            $this->pdo->rollBack();
            $this->pendingMessages = []; // 清除未发送的消息
            throw $e;
        }
    }
    private function flushMessages()
    {
        foreach ($this->pendingMessages as $message) {
            // 发送到 RabbitMQ / Kafka / Redis Stream 等
            $this->sendToMessageQueue($message);
        }
        $this->pendingMessages = [];
    }
}
// 使用示例
$messenger = new TransactionalMessenger($pdo);
$orderId = $messenger->executeInTransaction(function($tx) {
    // 执行数据库操作
    $orderId = insertOrder($tx->pdo, ...);
    // 注册事务提交后的消息
    $tx->sendAfterCommit('order.created', [
        'order_id' => $orderId,
        'user_id' => $user->id
    ]);
    return $orderId;
});

数据库触发器 + 消息消费者

使用数据库触发器和外部消费者的组合:

<?php
// MySQL 触发器示例
DELIMITER $$
CREATE TRIGGER after_order_insert
AFTER INSERT ON orders
FOR EACH ROW
BEGIN
    INSERT INTO message_outbox (message_type, payload, status)
    VALUES ('ORDER_CREATED', JSON_OBJECT('order_id', NEW.id), 'PENDING');
END$$
DELIMITER ;
// PHP 侧的消费者
class MessageWorker
{
    public function processOutbox()
    {
        $pdo = new PDO('mysql:host=localhost;dbname=app', 'user', 'pass');
        while (true) {
            // 查询待处理的消息
            $sql = "SELECT * FROM message_outbox 
                    WHERE status = 'PENDING' 
                    LIMIT 10 FOR UPDATE SKIP LOCKED";
            $messages = $pdo->query($sql)->fetchAll();
            foreach ($messages as $message) {
                try {
                    // 发送到消息队列
                    $this->mq->publish($message['message_type'], 
                                      json_decode($message['payload'], true));
                    // 标记发送成功
                    $pdo->exec("UPDATE message_outbox SET status = 'SENT' 
                               WHERE id = " . $message['id']);
                } catch (Exception $e) {
                    // 记录失败,稍后重试
                    $pdo->exec("UPDATE message_outbox SET status = 'FAILED' 
                               WHERE id = " . $message['id']);
                }
            }
            sleep(1); // 1秒轮询一次
        }
    }
}

使用可靠的 PHP 消息库

利用专门为 PHP 设计的分布式消息库:

<?php
// 使用 kafka-php 或 php-amqplib (RabbitMQ)
// 1. RabbitMQ 使用 Confirm 模式确保消息可靠性
use PhpAmqpLib\Connection\AMQPStreamConnection;
use PhpAmqpLib\Message\AMQPMessage;
class TransactionalPublisher
{
    private $connection;
    public function publishWithConfirmation($queue, $message)
    {
        $channel = $this->connection->channel();
        // 开启确认模式
        $channel->confirm_select();
        $msg = new AMQPMessage(json_encode($message), [
            'delivery_mode' => AMQPMessage::DELIVERY_MODE_PERSISTENT
        ]);
        $channel->basic_publish($msg, '', $queue);
        // 等待确认,超时则失败
        $channel->wait_for_pending_acks_returns(5);
        return true;
    }
}

最佳实践建议

完整的实现示例(基于 Redis Stream)

<?php
class TransactionalOrderService
{
    private $pdo;
    private $redis;
    private $streamKey = 'order_stream';
    public function __construct(PDO $pdo, Redis $redis)
    {
        $this->pdo = $pdo;
        $this->redis = $redis;
    }
    public function createOrderWithMessage($orderData)
    {
        $pdo = $this->pdo;
        $redis = $this->redis;
        // 1. 开始数据库事务
        $pdo->beginTransaction();
        try {
            // 2. 写入订单
            $orderId = $this->insertOrder($orderData);
            // 3. 写入本地消息表
            $messageId = $this->insertMessageRecord($orderId, $orderData);
            // 4. 提交事务
            $pdo->commit();
            // 5. 提交成功后发布到 Redis Stream
            //    这里使用 try-catch,失败会触发重试机制
            try {
                $redis->xAdd($this->streamKey, '*', [
                    'message_id' => $messageId,
                    'order_id' => $orderId,
                    'payload' => json_encode($orderData),
                    'timestamp' => time()
                ]);
                // 标记消息已添加
                $pdo->exec("UPDATE message_outbox SET status = 'PUBLISHED' 
                           WHERE id = $messageId");
            } catch (Exception $e) {
                // 消息发送失败,记录错误,由消费者补偿
                error_log("Message publish failed: " . $e->getMessage());
            }
            return $orderId;
        } catch (Exception $e) {
            $pdo->rollBack();
            throw $e;
        }
    }
    // 消费者端
    public function consumeOrderMessages()
    {
        $group = 'order-consumer-group';
        // 创建消费者组
        try {
            $this->redis->xGroup('CREATE', $this->streamKey, $group, 0, true);
        } catch (RedisException $e) {
            // group already exists
        }
        while (true) {
            // 读取消息
            $messages = $this->redis->xReadGroup(
                $group, 
                'consumer-' . getmypid(),
                [$this->streamKey => '>'],
                10, 
                1000
            );
            foreach ($messages as $stream => $items) {
                foreach ($items as $msgId => $data) {
                    $this->processMessage($data);
                    // 确认处理完成
                    $this->redis->xAck($this->streamKey, $group, [$msgId]);
                }
            }
        }
    }
}

补偿机制

<?php
// 事务消息的补偿处理
class MessageCompensation
{
    public function compensateFailedMessages()
    {
        // 查找超时未确认的消息
        $sql = "SELECT * FROM message_outbox 
                WHERE status = 'PENDING' 
                AND create_time < NOW() - INTERVAL 5 MINUTE";
        $messages = $this->pdo->query($sql)->fetchAll();
        foreach ($messages as $message) {
            // 检查订单状态
            $order = $this->getOrder($message['order_id']);
            if ($order && $order['status'] === 'COMPLETED') {
                // 订单已创建但消息未发送,进行补偿
                $this->sendMessage($message['message_data']);
                $this->pdo->exec("UPDATE message_outbox SET status = 'SENT' 
                                 WHERE id = " . $message['id']);
            } else {
                // 订单未完成,撤销消息
                $this->pdo->exec("UPDATE message_outbox SET status = 'CANCELLED' 
                                 WHERE id = " . $message['id']);
            }
        }
    }
}

PHP 事务消息的实现虽然没有像 Java 那样有专门的框架支持,但通过合理的设计仍然可以实现可靠的一致性保证,核心要点:

  1. 本地消息表是最通用且可靠的方案
  2. 消息确认机制保证消息不丢失
  3. 补偿机制处理各种失败场景
  4. 选择合适的消息队列(RabbitMQ、Kafka、Redis Stream 等)

根据项目的规模和需求,可以选择最适合的方案,对于中小型项目,使用 Redis Stream 或 RabbitMQ 配合本地消息表就足够了。

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