本文目录导读:

在 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 那样有专门的框架支持,但通过合理的设计仍然可以实现可靠的一致性保证,核心要点:
- 本地消息表是最通用且可靠的方案
- 消息确认机制保证消息不丢失
- 补偿机制处理各种失败场景
- 选择合适的消息队列(RabbitMQ、Kafka、Redis Stream 等)
根据项目的规模和需求,可以选择最适合的方案,对于中小型项目,使用 Redis Stream 或 RabbitMQ 配合本地消息表就足够了。