RocketMQ PHP SDK
更新時間 2024-09-06 01:04:30
最近更新時間: 2024-09-06 01:04:30
分享文章
說明分布式消息服務RocketMQ兼容了社區版 HTTP SDK,您可以使用社區版 HTTP SDK接入分布式消息服務RocketMQ。
前提條件:
-
在PHP安裝目錄下的composer.json文件中加入社區PHP SDK 依賴。
-
使用Composer安裝依賴。
composer install
發送普通消息
<?phprequire "vendor/autoload.php";use MQ\Model\TopicMessage;use MQ\MQClient;class NormalProducerExample{
private $client;
private $producer;
public function __construct()
{
$this->client = new MQClient(
// 填寫分布式消息服務RocketMQ控制臺HTTP接入點
"${HTTP_ENDPOINT}",
// 填寫AccessKey,在管理控制臺創建
"${ACCESS_KEY}",
// 填寫SecretKey 在管理控制臺創建
"${SECRET_KEY}"
);
// 所屬的 Topic
$topic = "${TOPIC}";
// Topic所屬實例ID,默認實例為空NULL
$instanceId = "${INSTANCE_ID}";
$this->producer = $this->client->getProducer($instanceId, $topic);
}
public function run()
{
try {
for ($i = 1; $i <= 4; $i++) {
$publishMessage = new TopicMessage(
"Hello RocketMQ" // 消息內容
);
// 設置屬性
$publishMessage->putProperty("a", $i);
// 設置消息KEY
$publishMessage->setMessageKey("MessageKey");
$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 NormalProducerExample();$instance->run();?>
消費普通消息
<?phprequire "vendor/autoload.php";use MQ\MQClient;class ConsumerExample{
private $client;
private $consumer;
public function __construct()
{
$this->client = new MQClient(
// 填寫分布式消息服務RocketMQ控制臺HTTP接入點
"${HTTP_ENDPOINT}",
// 填寫AccessKey,在管理控制臺創建
"${ACCESS_KEY}",
// 填寫SecretKey 在管理控制臺創建
"${SECRET_KEY}"
);
// 所屬的 Topic
$topic = "${TOPIC}";
// 您在控制臺創建的 Consumer ID(Group ID)
$groupId = "${GROUP_ID}";
// Topic所屬實例ID,默認實例為空NULL
$instanceId = "${INSTANCE_ID}";
$this->consumer = $this->client->getConsumer($instanceId, $topic, $groupId, "TagA");
}
public function run()
{
// 在當前線程循環消費消息,建議是多開個幾個線程并發消費消息
while (True) {
try {
// 長輪詢消費消息
// 長輪詢表示如果topic沒有消息則請求會在服務端掛住3s,3s內如果有消息可以消費則立即返回
$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);
try {
$this->ackMessages($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());
}
}
}
print "ack finish\n";
}
}
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());
}
}
}
}}$instance = new ConsumerExample();$instance->run();?>