云消息队列 RocketMQ 版提供类似XA或Open XA的分布式事务功能,通过云消息队列 RocketMQ 版事务消息,能达到分布式事务的最终一致。本文提供使用HTTP协议下的PHP SDK收发事务消息的示例代码。
背景信息
事务消息的交互流程如下图所示。
更多信息,请参见事务消息。
前提条件
您已完成以下操作:
安装PHP SDK。更多信息,请参见准备环境。
创建资源。代码中涉及的资源信息,例如实例、Topic和Group ID等,需要在控制台上提前创建。更多信息,请参见创建资源。
获取阿里云访问密钥AccessKey ID和AccessKey Secret。更多信息,请参见创建AccessKey。
发送事务消息
发送事务消息的示例代码如下。
<?php
require "vendor/autoload.php";
use MQ\Model\TopicMessage;
use MQ\MQClient;
class ProducerTest
{
private $client;
private $transProducer;
private $count;
private $popMsgCount;
public function __construct()
{
$this->client = new MQClient(
// 设置HTTP协议客户端接入点,进入消息队列RocketMQ版控制台实例详情页面的接入点区域查看。
"${HTTP_ENDPOINT}",
// 请确保环境变量ALIBABA_CLOUD_ACCESS_KEY_ID、ALIBABA_CLOUD_ACCESS_KEY_SECRET已设置。
// AccessKey ID,阿里云身份验证标识。
getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
// AccessKey Secret,阿里云身份验证密钥。
getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET')
);
// 消息所属的Topic,在消息队列RocketMQ版控制台创建。
$topic = "${TOPIC}";
// 您在消息队列RocketMQ版控制台创建的Group ID。
$groupId = "${GROUP_ID}";
// Topic所属的实例ID,在消息队列RocketMQ版控制台创建。
// 若实例有命名空间,则实例ID必须传入;若实例无命名空间,则实例ID传入null空值或字符串空值。实例的命名空间可以在消息队列RocketMQ版控制台的实例详情页面查看。
$instanceId = "${INSTANCE_ID}";
$this->transProducer = $this->client->getTransProducer($instanceId,$topic, $groupId);
$this->count = 0;
$this->popMsgCount = 0;
}
function processAckError($e) {
if ($e instanceof MQ\Exception\AckMessageException) {
// 如果Commit或Rollback时超过了TransCheckImmunityTime(针对发送事务消息的句柄)或者超过NextConsumeTime(针对consumeHalfMessage的句柄),则Commit或Rollback会失败。
printf("Commit/Rollback Error, RequestId:%s\n", $e->getRequestId());
foreach ($e->getAckMessageErrorItems() as $errorItem) {
printf("\tReceiptHandle:%s, ErrorCode:%s, ErrorMsg:%s\n", $errorItem->getReceiptHandle(), $errorItem->getErrorCode(), $errorItem->getErrorCode());
}
} else {
print_r($e);
}
}
function consumeHalfMsg() {
while($this->count < 3 && $this->popMsgCount < 15) {
$this->popMsgCount++;
try {
$messages = $this->transProducer->consumeHalfMessage(4, 3);
} catch (\Exception $e) {
if ($e instanceof MQ\Exception\MessageNotExistException) {
print "no half transaction message\n";
continue;
}
print_r($e->getMessage() . "\n");
sleep(3);
continue;
}
foreach ($messages as $message) {
printf("ID:%s TAG:%s BODY:%s \nPublishTime:%d, FirstConsumeTime:%d\nConsumedTimes:%d, NextConsumeTime:%d\nPropA:%s\n",
$message->getMessageId(), $message->getMessageTag(), $message->getMessageBody(),
$message->getPublishTime(), $message->getFirstConsumeTime(), $message->getConsumedTimes(), $message->getNextConsumeTime(),
$message->getProperty("a"));
print_r($message->getProperties());
$propA = $message->getProperty("a");
$consumeTimes = $message->getConsumedTimes();
try {
if ($propA == "1") {
print "\n commit transaction msg: " . $message->getMessageId() . "\n";
$this->transProducer->commit($message->getReceiptHandle());
$this->count++;
} else if ($propA == "2" && $consumeTimes > 1) {
print "\n commit transaction msg: " . $message->getMessageId() . "\n";
$this->transProducer->commit($message->getReceiptHandle());
$this->count++;
} else if ($propA == "3") {
print "\n rollback transaction msg: " . $message->getMessageId() . "\n";
$this->transProducer->rollback($message->getReceiptHandle());
$this->count++;
} else {
print "\n unknown transaction msg: " . $message->getMessageId() . "\n";
}
} catch (\Exception $e) {
$this->processAckError($e);
}
}
}
}
public function run()
{
// 循环发送4条事务消息。
for ($i = 0; $i < 4; $i++) {
$pubMsg = new TopicMessage("hello,mq");
// 设置消息的自定义属性。
$pubMsg->putProperty("a", $i);
// 设置消息的Key。
$pubMsg->setMessageKey("MessageKey");
// 设置事务第一次回查的时间,为相对时间。单位:秒,范围:10~300。
// 第一次事务回查后如果消息没有Commit或者Rollback,则之后每隔10s左右会回查一次,共回查24小时。
$pubMsg->setTransCheckImmunityTime(10);
$topicMessage = $this->transProducer->publishMessage($pubMsg);
print "\npublish -> \n\t" . $topicMessage->getMessageId() . " " . $topicMessage->getReceiptHandle() . "\n";
if ($i == 0) {
try {
// 发送完事务消息后能获取到半消息句柄,可以直接Commit或Rollback事务消息。
$this->transProducer->commit($topicMessage->getReceiptHandle());
print "\n commit transaction msg when publish: " . $topicMessage->getMessageId() . "\n";
} catch (\Exception $e) {
// 如果Commit或Rollback时超过了TransCheckImmunityTime则会失败。
$this->processAckError($e);
}
}
}
// 客户端需要有一个线程或者进程来消费没有确认的事务消息。
// 检查没有确认的事务消息。
$this->consumeHalfMsg();
}
}
$instance = new ProducerTest();
$instance->run();
?>
订阅事务消息
订阅事务消息的示例代码如下。
<?php
use MQ\MQClient;
require "vendor/autoload.php";
class ConsumerTest
{
private $client;
private $consumer;
public function __construct()
{
$this->client = new MQClient(
// 设置HTTP协议客户端接入点,进入消息队列RocketMQ版控制台实例详情页面的接入点区域查看。
"${HTTP_ENDPOINT}",
// 请确保环境变量ALIBABA_CLOUD_ACCESS_KEY_ID、ALIBABA_CLOUD_ACCESS_KEY_SECRET已设置。
// AccessKey ID,阿里云身份验证标识。
getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
// AccessKey Secret,阿里云身份验证密钥。
getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET')
);
// 消息所属的Topic,在消息队列RocketMQ版控制台创建。
$topic = "${TOPIC}";
// 您在消息队列RocketMQ版控制台创建的Group ID。
$groupId = "${GROUP_ID}";
// Topic所属的实例ID,在消息队列RocketMQ版控制台创建。
// 若实例有命名空间,则实例ID必须传入;若实例无命名空间,则实例ID传入null空值或字符串空值。实例的命名空间可以在消息队列RocketMQ版控制台的实例详情页面查看。
$instanceId = "${INSTANCE_ID}";
$this->consumer = $this->client->getConsumer($instanceId, $topic, $groupId);
}
public function ackMessages($receiptHandles)
{
try {
$this->consumer->ackMessage($receiptHandles);
} catch (\Exception $e) {
if ($e instanceof MQ\Exception\AckMessageException) {
// 某些消息的句柄可能超时,会导致消费确认失败。
printf("Ack Error, RequestId:%s\n", $e->getRequestId());
foreach ($e->getAckMessageErrorItems() as $errorItem) {
printf("\tReceiptHandle:%s, ErrorCode:%s, ErrorMsg:%s\n", $errorItem->getReceiptHandle(), $errorItem->getErrorCode(), $errorItem->getErrorCode());
}
}
}
}
public function run()
{
// 在当前线程循环消费消息,建议多开个几个线程并发消费消息。
while (True) {
try {
// 长轮询消费消息。
// 若Topic内没有消息,请求会在服务端挂起一段时间(长轮询时间),期间如果有消息可以消费则立即返回客户端。
$messages = $this->consumer->consumeMessage(
3, // 一次最多消费3条(最多可设置为16条)。
3 // 长轮询时间3秒(最多可设置为30秒)。
);
} catch (\MQ\Exception\MessageResolveException $e) {
// 当出现消息Body存在不合法字符,无法解析的时候,会抛出此异常。
// 可以正常解析的消息列表。
$messages = $e->getPartialResult()->getMessages();
// 无法正常解析的消息列表。
$failMessages = $e->getPartialResult()->getFailResolveMessages();
$receiptHandles = array();
foreach ($messages as $message) {
// 处理业务逻辑。
$receiptHandles[] = $message->getReceiptHandle();
printf("MsgID %s\n", $message->getMessageId());
}
foreach ($failMessages as $failMessage) {
// 处理存在不合法字符,无法解析的消息。
$receiptHandles[] = $failMessage->getReceiptHandle();
printf("Fail To Resolve Message. MsgID %s\n", $failMessage->getMessageId());
}
$this->ackMessages($receiptHandles);
continue;
} catch (\Exception $e) {
if ($e instanceof MQ\Exception\MessageNotExistException) {
// 没有消息可以消费,继续轮询。
printf("No message, contine long polling!RequestId:%s\n", $e->getRequestId());
continue;
}
print_r($e->getMessage() . "\n");
sleep(3);
continue;
}
print "consume finish, messages:\n";
// 处理业务逻辑。
$receiptHandles = array();
foreach ($messages as $message) {
$receiptHandles[] = $message->getReceiptHandle();
printf("MessageID:%s TAG:%s BODY:%s \nPublishTime:%d, FirstConsumeTime:%d, \nConsumedTimes:%d, NextConsumeTime:%d,MessageKey:%s\n",
$message->getMessageId(), $message->getMessageTag(), $message->getMessageBody(),
$message->getPublishTime(), $message->getFirstConsumeTime(), $message->getConsumedTimes(), $message->getNextConsumeTime(),
$message->getMessageKey());
print_r($message->getProperties());
}
// $message->getNextConsumeTime()前若不确认消息消费成功,则消息会被重复消费。
// 消息句柄有时间戳,同一条消息每次消费拿到的都不一样。
print_r($receiptHandles);
$this->ackMessages($receiptHandles);
print "ack finish\n";
}
}
}
$instance = new ConsumerTest();
$instance->run();
?>