O ApsaraMQ for RocketMQ entrega e consome mensagens ordenadas seguindo rigorosamente a ordem FIFO (first-in, first-out). O código de exemplo a seguir demonstra como enviar e consumir mensagens ordenadas com o SDK cliente HTTP para Go.
Tipos de ordenação
As mensagens ordenadas dividem-se em duas categorias:
|
Tipo |
Comportamento |
Quando usar |
|
Ordenação global |
Todas as mensagens de um tópico seguem uma única sequência FIFO. |
Quando cada mensagem precisa ser processada exatamente na ordem de envio, como na execução sequencial de comandos. |
|
Ordenação por partição |
As mensagens são distribuídas entre partições com base na chave de sharding. Cada partição mantém sua própria ordem FIFO. |
Recomendado quando apenas mensagens que compartilham um agrupamento lógico, como o mesmo ID de pedido ou ID de usuário, exigem ordenação estrita. |
Uma chave de sharding identifica a qual partição uma mensagem pertence. Ela difere da chave da mensagem: a chave da mensagem é um identificador de negócios usado para rastreamento, enquanto a chave de sharding controla a ordenação.
Para mais informações, consulte Mensagens ordenadas.
Como funciona a ordenação
A ordenação depende da coordenação entre o produtor e o consumidor.
Lado do envio: Utilize um único produtor e uma única thread para enviar mensagens. O broker preserva a sequência exata em que recebe as mensagens desse produtor.
Lado do consumo: Chame ConsumeMessageOrderly para puxar mensagens partição por partição. Dentro de cada partição, o próximo lote é entregue somente após o reconhecimento de todas as mensagens do lote atual. Caso o broker não receba um reconhecimento antes do prazo limite de NextConsumeTime, ele reentrega a mensagem não reconhecida.
Pré-requisitos
Conclua as configurações abaixo antes de executar o código de exemplo:
Crie uma instância, um tópico e um grupo de consumidores no console do ApsaraMQ for RocketMQ
Crie um par de AccessKey para sua conta Alibaba Cloud
Enviar mensagens ordenadas
O broker determina a ordem das mensagens pela sequência em que um único produtor ou thread as envia. Se múltiplos produtores ou threads enviarem simultaneamente, a ordem de chegada ao broker pode divergir da ordem de negócios pretendida. Para garantir a ordenação, envie a partir de um único produtor usando uma única thread.
O código a seguir envia oito mensagens com ordenação por partição. A chave de sharding i % 2 distribui as mensagens entre duas partições; assim, mensagens com índices pares vão para uma partição e índices ímpares para outra.
Substitua os espaços reservados antes de executar o código:
|
Espaço reservado |
Descrição |
Exemplo |
|
|
Endpoint HTTP obtido na página Instance Details no console do ApsaraMQ for RocketMQ |
|
|
|
Tópico criado no console |
|
|
|
ID da instância. Se a instância possuir um namespace, especifique o ID. Caso contrário, defina como uma string vazia. Verifique a página Instance Details para obter informações sobre o namespace. |
|
package main
import (
"fmt"
"time"
"strconv"
"os"
"github.com/aliyunmq/mq-http-go-sdk"
)
func main() {
// HTTP endpoint from the Instance Details page in the ApsaraMQ for RocketMQ console.
endpoint := "${HTTP_ENDPOINT}"
// Set the ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
// environment variables before running this code.
accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
// Topic created in the ApsaraMQ for RocketMQ console.
topic := "${TOPIC}"
// Instance ID. If the instance has a namespace, specify the ID.
// If not, set this to an empty string.
instanceId := "${INSTANCE_ID}"
client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")
mqProducer := client.GetProducer(instanceId, topic)
// Send 8 messages across 2 partitions.
for i := 0; i < 8; i++ {
msg := mq_http_sdk.PublishMessageRequest{
MessageBody: "hello mq!",
MessageTag: "",
Properties: map[string]string{},
}
msg.MessageKey = "MessageKey"
msg.Properties["a"] = strconv.Itoa(i)
// Sharding key determines the partition.
// Messages with the same sharding key are delivered in FIFO order.
msg.ShardingKey = strconv.Itoa(i % 2)
ret, err := mqProducer.PublishMessage(msg)
if err != nil {
fmt.Println(err)
return
} else {
fmt.Printf("Publish ---->\n\tMessageId:%s, BodyMD5:%s, \n", ret.MessageId, ret.MessageBodyMD5)
}
time.Sleep(time.Duration(100) * time.Millisecond)
}
}
Consumir mensagens ordenadas
Este código consome mensagens ordenadas usando long polling. O consumidor puxa mensagens partição por partição e reconhece cada lote antes de solicitar o próximo.
Substitua os espaços reservados pelos valores reais. Além dos espaços reservados listados na seção Enviar mensagens ordenadas, especifique o seguinte:
|
Espaço reservado |
Descrição |
Exemplo |
|
|
Grupo de consumidores criado no console do ApsaraMQ for RocketMQ |
|
package main
import (
"fmt"
"github.com/gogap/errors"
"strings"
"time"
"os"
"github.com/aliyunmq/mq-http-go-sdk"
)
func main() {
// HTTP endpoint from the Instance Details page in the ApsaraMQ for RocketMQ console.
endpoint := "${HTTP_ENDPOINT}"
// Set the ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
// environment variables before running this code.
accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
secretKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
// Topic created in the ApsaraMQ for RocketMQ console.
topic := "${TOPIC}"
// Instance ID. If the instance has a namespace, specify the ID.
// If not, set this to an empty string.
instanceId := "${INSTANCE_ID}"
// Consumer group created in the ApsaraMQ for RocketMQ console.
groupId := "${GROUP_ID}"
client := mq_http_sdk.NewAliyunMQClient(endpoint, accessKey, secretKey, "")
mqConsumer := client.GetConsumer(instanceId, topic, groupId, "")
for {
endChan := make(chan int)
respChan := make(chan mq_http_sdk.ConsumeMessageResponse)
errChan := make(chan error)
go func() {
select {
case resp := <-respChan:
{
var handles []string
fmt.Printf("Consume %d messages---->\n", len(resp.Messages))
for _, v := range resp.Messages {
handles = append(handles, v.ReceiptHandle)
fmt.Printf("\tMessageID: %s, PublishTime: %d, MessageTag: %s\n"+
"\tConsumedTimes: %d, FirstConsumeTime: %d, NextConsumeTime: %d\n"+
"\tBody: %s\n"+
"\tProps: %s\n"+
"\tShardingKey: %s\n",
v.MessageId, v.PublishTime, v.MessageTag, v.ConsumedTimes,
v.FirstConsumeTime, v.NextConsumeTime, v.MessageBody, v.Properties, v.ShardingKey)
}
// Acknowledge all messages in this batch.
// If the broker does not receive an ACK before NextConsumeTime,
// it redelivers the message.
ackerr := mqConsumer.AckMessage(handles)
if ackerr != nil {
fmt.Println(ackerr)
if errAckItems, ok := ackerr.(errors.ErrCode).Context()["Detail"].([]mq_http_sdk.ErrAckItem); ok {
for _, errAckItem := range errAckItems {
fmt.Printf("\tErrorHandle:%s, ErrorCode:%s, ErrorMsg:%s\n",
errAckItem.ErrorHandle, errAckItem.ErrorCode, errAckItem.ErrorMsg)
}
} else {
fmt.Println("ack err =", ackerr)
}
time.Sleep(time.Duration(3) * time.Second)
} else {
fmt.Printf("Ack ---->\n\t%s\n", handles)
}
endChan <- 1
}
case err := <-errChan:
{
if strings.Contains(err.(errors.ErrCode).Error(), "MessageNotExist") {
fmt.Println("\nNo new message, continue!")
} else {
fmt.Println(err)
time.Sleep(time.Duration(3) * time.Second)
}
endChan <- 1
}
case <-time.After(35 * time.Second):
{
fmt.Println("Timeout of consumer message ??")
endChan <- 1
}
}
}()
// The consumer pulls partitionally ordered messages partition by partition.
// Within each partition, the next batch is delivered only after all messages
// in the current batch are acknowledged.
// Long polling: if no message is available, the request is held on the broker
// for up to the specified polling duration. The network timeout is 35 seconds.
mqConsumer.ConsumeMessageOrderly(respChan, errChan,
3, // Max messages per batch (up to 16).
3, // Long polling duration in seconds (up to 30).
)
<-endChan
}
}
Principais comportamentos de consumo
|
Comportamento |
Detalhe |
|
Ordenação no nível da partição |
O consumidor pode puxar de múltiplas partições simultaneamente, mas as mensagens dentro de cada partição são sempre entregues na ordem de envio. |
|
Barreira de reconhecimento em lote |
O próximo lote de uma partição não é entregue até que todas as mensagens do lote anterior sejam reconhecidas. |
|
Reentrega automática |
Se o broker não receber nenhum reconhecimento antes de |
|
Long polling |
Quando não há mensagens disponíveis, o broker retém a requisição por até a duração de polling especificada (3 segundos neste exemplo) antes de retornar uma resposta vazia. O timeout de rede é de 35 segundos. |
Próximos passos
Mensagens ordenadas: Tipos de ordenação, garantias e considerações de design
Envie e receba outros tipos de mensagens usando o SDK HTTP para Go
Prepare o ambiente: Configure o ambiente de desenvolvimento Go