全部產品
Search
文件中心

ApsaraMQ for RocketMQ:收發事務訊息

更新時間:Jul 01, 2024

雲訊息佇列 RocketMQ 版提供類似XA或Open XA的分散式交易功能,通過雲訊息佇列 RocketMQ 版事務訊息,能達到分散式交易的最終一致。本文提供使用HTTP協議下的PHP SDK收發事務訊息的範例程式碼。

背景資訊

事務訊息的互動流程如下圖所示。

圖片1.png

更多資訊,請參見事務訊息

前提條件

您已完成以下操作:

  • 安裝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();


?>