全部产品
Search
文档中心

云消息队列 RocketMQ 版:收发顺序消息

更新时间:Jul 25, 2023

顺序消息(FIFO消息)是云消息队列 RocketMQ 版提供的一种严格按照顺序来发布和消费的消息类型。本文提供使用HTTP协议下的PHP SDK收发顺序消息的示例代码。

背景信息

顺序消息分为两类:

  • 全局顺序:对于指定的一个Topic,所有消息按照严格的先入先出FIFO(First In First Out)的顺序进行发布和消费。

  • 分区顺序:对于指定的一个Topic,所有消息根据Sharding Key进行区块分区。同一个分区内的消息按照严格的FIFO顺序进行发布和消费。Sharding Key是顺序消息中用来区分不同分区的关键字段,和普通消息的Key是完全不同的概念。

更多信息,请参见顺序消息

前提条件

您已完成以下操作:

  • 安装PHP SDK。更多信息,请参见准备环境

  • 创建资源。代码中涉及的资源信息,例如实例、Topic和Group ID等,需要在控制台上提前创建。更多信息,请参见创建资源

  • 获取阿里云访问密钥AccessKey ID和AccessKey Secret。更多信息,请参见创建AccessKey

发送顺序消息

重要

云消息队列 RocketMQ 版服务端判定消息产生的顺序性是参照单一生产者、单一线程并发下消息发送的时序。如果发送方有多个生产者或者有多个线程并发发送消息,则此时只能以到达云消息队列 RocketMQ 版服务端的时序作为消息顺序的依据,和业务侧的发送顺序未必一致。

发送顺序消息的示例代码如下。

<?php

require "vendor/autoload.php";

use MQ\Model\TopicMessage;
use MQ\MQClient;

class ProducerTest
{
    private $client;
    private $producer;

    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}";
        // Topic所属的实例ID,在消息队列RocketMQ版控制台创建。
        // 若实例有命名空间,则实例ID必须传入;若实例无命名空间,则实例ID传入null空值或字符串空值。实例的命名空间可以在消息队列RocketMQ版控制台的实例详情页面查看。
        $instanceId = "${INSTANCE_ID}";

        $this->producer = $this->client->getProducer($instanceId, $topic);
    }

    public function run()
    {
        try
        {
            for ($i=1; $i<=4; $i++)
            {
                $publishMessage = new TopicMessage(
                    "hello mq!"// 消息内容。
                );
                // 设置消息的自定义属性。
                $publishMessage->putProperty("a", $i);
                // 设置分区顺序消息的Sharding Key,用于标识不同的分区。Sharding Key与消息的Key是完全不同的概念。
                $publishMessage->setShardingKey($i % 2);

                $result = $this->producer->publishMessage($publishMessage);

                print "Send mq message success. msgId is:" . $result->getMessageId() . ", bodyMD5 is:" . $result->getMessageBodyMD5() . "\n";
            }
        } catch (\Exception $e) {
            print_r($e->getMessage() . "\n");
        }
    }
}


$instance = new ProducerTest();
$instance->run();

?>

            

订阅顺序消息

订阅顺序消息的示例代码如下。

<?php

require "vendor/autoload.php";

use MQ\MQClient;

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->consumeMessageOrderly(
                    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,ShardingKey:%s\n",
                    $message->getMessageId(), $message->getMessageTag(), $message->getMessageBody(),
                    $message->getPublishTime(), $message->getFirstConsumeTime(), $message->getConsumedTimes(), $message->getNextConsumeTime(),
                    $message->getShardingKey());
                print_r($message->getProperties());
            }

            // $message->getNextConsumeTime()前若不确认消息消费成功,则消息会被重复消费。
            // 消息句柄有时间戳,同一条消息每次消费拿到的都不一样。
            print_r($receiptHandles);
            $this->ackMessages($receiptHandles);
            print "=======>ack finish\n";


        }

    }
}


$instance = new ConsumerTest();
$instance->run();

?>