O ApsaraMQ for RocketMQ entrega mensagens ordenadas seguindo rigorosamente a ordem FIFO (first-in-first-out). Este tópico fornece códigos de exemplo em PHP para enviar e receber mensagens ordenadas por meio do SDK cliente HTTP.
Tipos de ordenação
O ApsaraMQ for RocketMQ oferece suporte a dois escopos de ordenação:
|
Tipo |
Escopo de ordenação |
Funcionamento |
|
Ordenação global |
Todo o tópico |
Todas as mensagens do tópico são enviadas e consumidas em ordem FIFO |
|
Ordenação por partição |
Por partição |
As mensagens são distribuídas nas partições com base na chave de sharding. Dentro de cada partição, o consumo ocorre em ordem FIFO |
A chave de sharding identifica a qual partição uma mensagem pertence. Ela difere da chave da mensagem.
Para mais informações, consulte Mensagens ordenadas.
Garantias de ordenação
A entrega ordenada depende do comportamento correto tanto do produtor quanto do consumidor.
O broker determina a ordem das mensagens com base na sequência em que um único produtor ou thread as envia. Se vários produtores ou threads enviarem mensagens simultaneamente, a ordem será definida pela chegada ao broker, o que pode divergir da ordem lógica de negócios pretendida.
Lado do consumidor:
O consumidor pode buscar mensagens ordenadas de várias partições em um único lote. As mensagens de cada partição mantêm a ordem de envio original.
Confirme (ACK) todas as mensagens de uma partição antes de receber o próximo lote dessa mesma partição.
Caso o broker não receba um ACK antes do tempo limite definido por
getNextConsumeTime(), a mensagem será reentregue.
Pré-requisitos
Antes de começar, verifique se você possui:
SDK PHP para ApsaraMQ for RocketMQ instalado. Para detalhes, consulte Preparar o ambiente
Instância, tópico e grupo de consumidores do ApsaraMQ for RocketMQ criados no console do ApsaraMQ for RocketMQ
Par de AccessKey para sua conta Alibaba Cloud. Para detalhes, consulte Criar um par de AccessKey
Enviar mensagens ordenadas
O código abaixo envia quatro mensagens ordenadas e as distribui em duas partições usando setShardingKey(). Mensagens que compartilham a mesma chave de sharding são entregues em ordem FIFO dentro dessa partição.
<?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 endpoint. Find this in the HTTP Endpoint section on the Instance Details
// page in the ApsaraMQ for RocketMQ console.
"${HTTP_ENDPOINT}",
// Get AccessKey credentials from environment variables.
getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET')
);
// Topic created in the ApsaraMQ for RocketMQ console.
$topic = "${TOPIC}";
// Instance ID. If the instance has no namespace, set this to null or "".
// Check the Instance Details page for namespace information.
$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! " // Message body
);
// Set a custom message property.
$publishMessage->putProperty("a", $i);
// Set the sharding key to distribute messages across partitions.
// Messages with the same sharding key are delivered in FIFO order.
$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();
?>
Substitua os placeholders a seguir pelos valores reais:
|
Placeholder |
Descrição |
|
|
Endpoint HTTP obtido na página de detalhes da instância |
|
|
Nome do tópico |
|
|
ID da instância (null ou "" se não houver namespace) |
Receber mensagens ordenadas
O código a seguir consome mensagens ordenadas utilizando consumeMessageOrderly(), preservando a ordem FIFO dentro de cada partição. O consumidor emprega long polling: se nenhuma mensagem estiver disponível, a solicitação permanece no broker até a chegada de uma nova mensagem ou o esgotamento do tempo limite de polling.
<?php
require "vendor/autoload.php";
use MQ\MQClient;
class ConsumerTest
{
private $client;
private $consumer;
public function __construct()
{
$this->client = new MQClient(
// HTTP endpoint. Find this in the HTTP Endpoint section on the Instance Details
// page in the ApsaraMQ for RocketMQ console.
"${HTTP_ENDPOINT}",
// Get AccessKey credentials from environment variables.
getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET')
);
// Topic created in the ApsaraMQ for RocketMQ console.
$topic = "${TOPIC}";
// Consumer group ID created in the ApsaraMQ for RocketMQ console.
$groupId = "${GROUP_ID}";
// Instance ID. If the instance has no namespace, set this to null or "".
// Check the Instance Details page for namespace information.
$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) {
// Handle ACK failures. This occurs when a receipt handle has expired.
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()
{
// Poll for messages continuously. For production use, run multiple threads
// to consume messages concurrently across partitions.
while (True) {
try {
// consumeMessageOrderly() ensures FIFO order within each partition.
// If the broker does not receive an ACK for a message, it redelivers
// that message before sending subsequent messages from the same partition.
$messages = $this->consumer->consumeMessageOrderly(
3, // Max messages per batch (up to 16)
3 // Long polling timeout in seconds (up to 30)
);
} catch (\MQ\Exception\MessageResolveException $e) {
// Handle messages that cannot be parsed due to invalid characters.
$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) {
// No messages available. Long polling continues automatically.
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());
}
// ACK all messages. If the broker does not receive an ACK before
// getNextConsumeTime(), it redelivers the message. Each redelivery
// assigns a new receipt handle with an updated timestamp.
print_r($receiptHandles);
$this->ackMessages($receiptHandles);
print "=======>ack finish\n";
}
}
}
$instance = new ConsumerTest();
$instance->run();
?>
Além dos placeholders mencionados na seção Enviar mensagens ordenadas, substitua o seguinte:
|
Placeholder |
Descrição |
|
|
ID do grupo de consumidores |
Parâmetros de consumeMessageOrderly()
|
Parâmetro |
Descrição |
Intervalo |
|
|
Quantidade máxima de mensagens a serem consumidas por lote |
Até 16 |
|
|
Duração do long polling em segundos. Na ausência de mensagens, a requisição fica retida no broker até a chegada de uma mensagem ou o fim do tempo limite |
Até 30 |
Tratamento de erros
O consumidor lida com três cenários de erro:
|
Exceção |
Causa |
Ação |
|
|
Caracteres inválidos no corpo da mensagem impedem a análise sintática |
Confirmar (ACK) tanto as mensagens analisáveis quanto as não analisáveis |
|
|
Nenhuma mensagem disponível no tópico |
Continuar o long polling automaticamente |
|
|
O receipt handle expirou antes do envio do ACK |
Registrar o erro em log e tentar novamente na próxima entrega |